Lab 06: Lakeflow & DLT (Bronze to Gold)
Before we write code, let's understand our Goal: We need to ingest dirty, raw data, clean it up, and aggregate it for the business, ensuring that bad data never makes it to the final dashboard.
The Tool: Delta Live Tables (DLT) is Databricks' declarative framework. Instead of writing complex step-by-step logic, you just define what you want the final tables to look like, and DLT handles the underlying dependencies and infrastructure.
1. The Medallion Architecture
Real-World Analogy Mapping: Imagine a water filtration plant.
- Bronze (The River): The raw, muddy water exactly as it flowed in from nature. We never drink this, but we store it just in case we need to filter it again later.
- Silver (The Filtered Tank): The water is purified. Impurities (nulls, duplicates) are removed.
- Gold (The Bottled Water): The water is perfectly packaged and aggregated, ready for consumers (BI Dashboards) to drink instantly.
2. Full DLT Pipeline Implementation
Problems Faced: In traditional PySpark, if a Bronze table failed to update, the Silver table would still run, processing old data. Engineers had to use complex orchestration tools like Airflow to map dependencies.
How Present Technology Solves It:
With DLT, you just write Python functions decorated with @dlt.table. DLT automatically reads the code, figures out that Silver depends on Bronze, and builds the dependency graph for you.
import dlt
from pyspark.sql.functions import col
# 1. BRONZE: The River (Raw Ingestion)
@dlt.table(name="bronze_orders")
def ingest_bronze():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.load("s3://raw-data/orders/")
)
# 2. SILVER: The Filtered Tank (Cleansed Data)
# Note: We use "Expectations" to act as the filter. Bad rows are dropped automatically!
@dlt.table(name="silver_orders")
@dlt.expect_or_drop("valid_order", "order_amount > 0")
@dlt.expect_or_fail("valid_id", "order_id IS NOT NULL")
def clean_silver():
return (
dlt.read_stream("bronze_orders")
.dropDuplicates(["order_id"]) # Deduplication
.withColumn("order_amount", col("order_amount").cast("double"))
)
# 3. GOLD: The Bottled Water (Business Aggregations)
# We use a Materialized View (read) instead of a Stream (read_stream) because streaming aggregations are complex.
@dlt.table(name="gold_daily_revenue")
def aggregate_gold():
return (
dlt.read("silver_orders")
.groupBy("order_date")
.agg({"order_amount": "sum"})
)
3. DLT Event Logs and Data Quality Monitoring
Follow the Data:
When a row is dropped because the order_amount was negative, where does that information go?
DLT automatically logs every single event (rows processed, rows dropped, expectation failures) to a hidden Delta table called the Event Log. You can query this log to build data-quality monitoring dashboards to see exactly why your data was filtered out.
SELECT timestamp, details
FROM delta.`s3://metadata/dlt_pipeline_storage/system/events`
WHERE event_type = 'flow_progress';
← Previous: Lab 05: CDC and SCD Implementation | Next: Lab 07: Unity Catalog & Governance SQL →**