Example source: runners/databricks/replay_spark.ipynb¶
Source revision: 75f65139e26eb7079b7a897c877de97859c0120b
This page is the generated source; it shows the complete readable projection of the canonical file.
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 Replay¶
Attach verified dependencies before execution. Range semantics are [START, END).
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
START = "2025-01-01"
END = "2025-02-01"
CHUNK_INTERVAL_JSON = "null"
CHUNK_COLUMN = ""
SAVE_WATERMARK = "false"
CONFIRM_SAVE_WATERMARK = "false"
import json
import re
from importlib.metadata import PackageNotFoundError, version
def widget_value(name, default):
try:
return dbutils.widgets.get(name) # type: ignore[name-defined]
except Exception:
try:
dbutils.widgets.text(name, str(default)) # type: ignore[name-defined]
return dbutils.widgets.get(name) # type: ignore[name-defined]
except Exception as exc:
raise RuntimeError(f"Unable to read or create Databricks widget {name}") from exc
def parse_bool(value):
normalized = value.strip().lower()
if normalized not in {"true", "false"}:
raise ValueError("boolean widget must be true or false")
return normalized == "true"
def decode_boundary(value):
return int(value) if re.fullmatch(r"[+-]?\d+", value) else value
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))
start = decode_boundary(widget_value("START", START))
end = decode_boundary(widget_value("END", END))
chunk_interval = json.loads(widget_value("CHUNK_INTERVAL_JSON", CHUNK_INTERVAL_JSON))
chunk_column = widget_value("CHUNK_COLUMN", CHUNK_COLUMN).strip() or None
save_watermark = parse_bool(widget_value("SAVE_WATERMARK", SAVE_WATERMARK))
confirm_save_watermark = parse_bool(widget_value("CONFIRM_SAVE_WATERMARK", CONFIRM_SAVE_WATERMARK))
if save_watermark and not confirm_save_watermark:
raise ValueError("SAVE_WATERMARK requires CONFIRM_SAVE_WATERMARK")
from datacoolie.core.models.run_config import DataCoolieRunConfig, ReplayConfig
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)
replay = ReplayConfig(start=start, end=end, chunk_interval=chunk_interval, save_watermark=save_watermark, chunk_column=chunk_column)
failed = 0
with DataCoolieDriver(engine=engine, metadata_provider=metadata, log_base_path=log_base_path, config=DataCoolieRunConfig(job_num=job_num, job_index=job_index, stop_on_error=True, allowed_function_prefixes=[])) as driver:
dataflows = driver.load_dataflows(stage=stage)
result = driver.run_replay(dataflows=dataflows, replay=replay)
failed = result.failed
if failed:
raise RuntimeError(f"DataCoolie replay failed for {failed} dataflows")