Lab 05: CDC and SCD Implementation
Before we write code, let's understand our Goal: Imagine a customer named Alice lives in "New York". Tomorrow, she moves to "Boston". Do we delete "New York" and replace it with "Boston"? Or do we keep a historical record that she used to live in New York?
This is called a Slowly Changing Dimension (SCD).
1. SCD Type 1 (Overwrite History)
Goal: Overwrite the old record. There is no historical tracking.
Anatomy Breakdown:
If Alice moves to Boston, the MERGE command looks at her customer_id. Since the ID exists in the target table (MATCHED), it updates her address. If a brand new customer signs up (NOT MATCHED), it inserts them.
MERGE INTO silver.customers target
USING bronze.customer_updates source
ON target.customer_id = source.customer_id
WHEN MATCHED THEN
UPDATE SET target.address = source.address
WHEN NOT MATCHED THEN
INSERT (customer_id, address) VALUES (source.customer_id, source.address)
2. SCD Type 2 (Preserve History)
Goal: Preserve history by expiring the old record (setting is_current = false and setting an end_date) and inserting the new record as the current active version.
What Existed Previously:
Writing SCD Type 2 logic in standard PySpark or SQL requires hundreds of lines of complex LEFT JOINs, timestamp comparisons, and window functions. It is incredibly easy to make a mistake and corrupt the historical timeline.
How Present Technology Solves It:
Delta Live Tables (DLT) natively solves this using a declarative function called APPLY CHANGES INTO. You just tell DLT which columns to track, and it automatically builds the is_current, start_date, and end_date columns for you!
import dlt
from pyspark.sql.functions import col, expr
# 1. Define the target table (This is where the history will be stored)
dlt.create_streaming_table("silver_customers_scd2")
# 2. The Smart Way: Apply the CDC logic declaratively
dlt.apply_changes(
target = "silver_customers_scd2",
source = "bronze_customer_cdc",
keys = ["customer_id"], # The unique ID to track
sequence_by = col("update_timestamp"), # Handles out-of-order events (e.g. Wednesday's update arrives before Tuesday's)
apply_as_deletes = expr("operation_type = 'DELETE'"), # Handles hard deletes from the source database
except_column_list = ["operation_type", "update_timestamp"],
stored_as_scd_type = 2, # MAGIC: Automatically creates is_current, start_date, and end_date columns!
track_history_column_list = ["address", "subscription_tier"] # Only create a new history row if these specific columns change
)
← Previous: Lab 04: Delta Lake Advanced Operations | Next: Lab 06: Lakeflow & DLT (Bronze to Gold) →**