A config-driven data ingestion and transformation framework built on Databricks Free Edition and Unity Catalog. One generic pipeline, driven entirely by a control table — onboarding a new source means adding a row to the config, not writing new code.
Most tutorials teach "one notebook per table." Production data engineering teams don't work that way — they build a single generic framework that reads a metadata/control table to decide how to load each source (format, keys, SCD strategy, target table), and apply the same logic to every entity. This project implements that pattern end-to-end: generic Auto Loader ingestion, generic SCD Type 1 / SCD Type 2 / append-only merge logic, a data-quality gate with quarantine, and full audit logging — all dispatched from one control table.
meta.ingestion_config (control table)
│
▼
01_data_generator.py ──► Unity Catalog Volume "landing"
│
▼
02_bronze_ingestion.py ──► Auto Loader, loops over config ──► Bronze Delta tables
│ │
│ audit.pipeline_run_log
│ audit.error_log
▼
03_silver_merge_scd.py ──► generic MERGE / SCD2 / append, driven by config ──► Silver Delta tables
▼
04_gold_aggregations.py ──► Customer 360, Daily Sales, Inventory Health, Pipeline Health
▼
Databricks SQL Dashboard
▲
05_workflow_job.json ──► Databricks Workflows (retries, timeouts, schedule)
- Databricks Free Edition — serverless compute, Unity Catalog, Volumes
- Delta Lake — MERGE, schema evolution, ACID transactions
- Auto Loader (
cloudFiles) — schema inference with automatic rescue-column handling for schema drift - PySpark / Delta MERGE — generic SCD Type 1, full SCD Type 2 (effective/end dates, hash-based change detection), and append-only fact loading
- Databricks Workflows — multi-task orchestration with retries, timeouts, and scheduling
- Metadata-driven design — the control table (
meta.ingestion_config) is the single source of truth for how each entity is ingested and transformed; the notebooks contain no hard-coded per-table logic - Generic SCD Type 1 and Type 2 implementations — dispatched dynamically by
load_type, including hash-based change detection for SCD2 rather than comparing every column by hand - Data quality as a first-class step — rows missing primary keys are quarantined with a reason code, not silently dropped or allowed to break downstream joins
- Full observability — every run (success or failure, per entity, per layer) is logged
to an audit table, which itself feeds a
gold_pipeline_healthtable for the dashboard - Resilient orchestration — one entity failing doesn't stop the others; errors are caught, logged, and the pipeline continues
metadata_etl_framework/
├── 00_setup_metadata_tables.py # Catalog/schema setup, control table, audit tables
├── 01_data_generator.py # Simulated source system (includes intentional dirty data)
├── 02_bronze_ingestion.py # Config-driven Auto Loader ingestion + audit logging
├── 03_silver_merge_scd.py # Generic SCD1 / SCD2 / append merge logic + DQ gate
├── 04_gold_aggregations.py # Business aggregates + pipeline health table
├── 05_workflow_job.json # Databricks Workflow: retries, timeouts, schedule
└── 06_dashboard_queries.sql # Queries powering the SQL dashboard
- Create a free workspace at databricks.com/try-databricks (choose Free Edition — compute is serverless-only, no cluster configuration needed).
- Import all files into a
metadata_etl_frameworkfolder in your Workspace. - Run
00_setup_metadata_tables.pyonce — creates the catalog, schemas, Volumes, control table, and audit tables. - Run
01_data_generator.py(mode:init). - Run
02_bronze_ingestion.py→03_silver_merge_scd.py→04_gold_aggregations.py, in order — or skip straight to step 6. - Import
05_workflow_job.jsonas a Databricks Workflow (Workflows → Jobs → Create → JSON view) to run all steps end-to-end on a schedule. - Build a Databricks SQL dashboard from
06_dashboard_queries.sql, on top ofgold_customer_360,gold_daily_sales_summary,gold_inventory_health, andgold_pipeline_health.
- Add a new source: add one row to
meta.ingestion_config— no notebook changes needed. - Add per-entity DQ rules: extend the config table with a
dq_rulescolumn and branch on it insideapply_dq_gate. - Make Bronze continuous: swap
trigger(availableNow=True)for a continuous trigger in02_bronze_ingestion.pyto turn Bronze into an always-on streaming job.
MIT



