2 min read

    Lab 05: CDC and SCD Implementation

    #databricks#cdc#scd#lab

    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.

    sql
    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!

    python
    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) →**