Skip to content

Example source: runners/databricks/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 Databricks Spark Runner

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

METADATA_PATH = "/Volumes/main/default/datacoolie_example/metadata/metadata.json"
CONNECTIONS_PATH = ""
SCHEMA_HINTS_PATH = ""
WATERMARK_BASE_PATH = "/Volumes/main/default/datacoolie_example/.runtime/watermarks"
LOG_BASE_PATH = "/Volumes/main/default/datacoolie_example/.runtime/logs"
STAGE = ""
JOB_NUM = 1
JOB_INDEX = 0

from importlib.metadata import PackageNotFoundError, version
def widget_value(name, default):
    try:
        return dbutils.widgets.get(name)  # type: ignore[name-defined]
    except Exception:
        dbutils.widgets.text(name, str(default))  # type: ignore[name-defined]
        return dbutils.widgets.get(name)  # type: ignore[name-defined]

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

metadata_path = widget_value("METADATA_PATH", METADATA_PATH)
connections_path = widget_value("CONNECTIONS_PATH", CONNECTIONS_PATH) or None
schema_hints_path = widget_value("SCHEMA_HINTS_PATH", SCHEMA_HINTS_PATH) or None
watermark_base_path = widget_value("WATERMARK_BASE_PATH", WATERMARK_BASE_PATH)
log_base_path = widget_value("LOG_BASE_PATH", LOG_BASE_PATH)
stage = widget_value("STAGE", STAGE)
job_num = int(widget_value("JOB_NUM", JOB_NUM))
job_index = int(widget_value("JOB_INDEX", JOB_INDEX))

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.databricks_platform import DatabricksPlatform

platform = DatabricksPlatform(runtime="databricks")
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=job_num, job_index=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")