Skip to content

Example source: runners/databricks/maintenance_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 Maintenance

Attach verified dependencies before confirmed execution.

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"
CONNECTION = ""
JOB_NUM = 1
JOB_INDEX = 0
DO_COMPACT = "true"
DO_CLEANUP = "true"
RETENTION_HOURS = "168"
CONFIRM_MAINTENANCE = "false"

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"

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)
connection = widget_value("CONNECTION", CONNECTION) or None
job_num = int(widget_value("JOB_NUM", JOB_NUM))
job_index = int(widget_value("JOB_INDEX", JOB_INDEX))
do_compact = parse_bool(widget_value("DO_COMPACT", DO_COMPACT))
do_cleanup = parse_bool(widget_value("DO_CLEANUP", DO_CLEANUP))
retention_hours = int(widget_value("RETENTION_HOURS", RETENTION_HOURS))
confirm_maintenance = parse_bool(widget_value("CONFIRM_MAINTENANCE", CONFIRM_MAINTENANCE))
if not do_compact and not do_cleanup:
    raise ValueError("maintenance must enable compact, cleanup, or both")
if not confirm_maintenance:
    raise ValueError("maintenance mutation requires CONFIRM_MAINTENANCE")

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, retention_hours=retention_hours, stop_on_error=True, allowed_function_prefixes=[])

with DataCoolieDriver(engine=engine, metadata_provider=metadata, log_base_path=log_base_path, config=config) as driver:
    result = driver.run_maintenance(connection=connection, do_compact=do_compact, do_cleanup=do_cleanup)
    if result.failed:
        raise RuntimeError(f"DataCoolie maintenance failed for {result.failed} dataflows")