Skip to content

Example source: runners/fabric/run_spark.ipynb

Source revision: 75f65139e26eb7079b7a897c877de97859c0120b

This page is the generated source; it shows the complete readable projection of the canonical file.

raw · View repository source

This page is a generated, non-executed projection of the notebook.

The raw .ipynb file is the canonical notebook source; execution counts and outputs are intentionally omitted.

DataCoolie Fabric Spark Runner

Attach the verified DataCoolie and functions artifacts to the environment before execution. STAGE is passed unchanged to one framework call.

METADATA_PATH = "abfss://your-workspace@onelake.dfs.fabric.microsoft.com/your-lakehouse.Lakehouse/Files/datacoolie-example/metadata/metadata.json"
CONNECTIONS_PATH = None
SCHEMA_HINTS_PATH = None
WATERMARK_BASE_PATH = "abfss://your-workspace@onelake.dfs.fabric.microsoft.com/your-lakehouse.Lakehouse/Files/datacoolie-example/.runtime/watermarks"
LOG_BASE_PATH = "abfss://your-workspace@onelake.dfs.fabric.microsoft.com/your-lakehouse.Lakehouse/Files/datacoolie-example/.runtime/logs"
STAGE = ""
JOB_NUM = 1
JOB_INDEX = 0

from importlib.metadata import PackageNotFoundError, version
try:
    print(f"DataCoolie version: {version('datacoolie')}")
except PackageNotFoundError as exc:
    raise RuntimeError("Attach the verified DataCoolie package before running this notebook") from exc

from datacoolie.core.models.run_config import DataCoolieRunConfig
from datacoolie.engines.spark_engine import SparkEngine
from datacoolie.metadata.file_provider import FileProvider
from datacoolie.orchestration.driver import DataCoolieDriver
from datacoolie.platforms.fabric_platform import FabricPlatform

platform = FabricPlatform(runtime="fabric")
engine = SparkEngine(spark_session=spark, platform=platform)  # type: ignore[name-defined]
metadata = FileProvider(config_path=METADATA_PATH, connections_path=CONNECTIONS_PATH, schema_hints_path=SCHEMA_HINTS_PATH, platform=platform, watermark_base_path=WATERMARK_BASE_PATH)
config = DataCoolieRunConfig(job_num=int(JOB_NUM), job_index=int(JOB_INDEX), max_workers=4, stop_on_error=True, allowed_function_prefixes=[])

failed = 0
with DataCoolieDriver(engine=engine, metadata_provider=metadata, log_base_path=LOG_BASE_PATH, config=config) as driver:
    result = driver.run(stage=STAGE)
    failed = result.failed
if failed:
    raise RuntimeError(f"DataCoolie run failed for {failed} dataflows")