October 24, 2024

Spark SQL MERGE on Delta Lake: SCD Type 2 for the Coverholder Dimension

Slowly Changing Dimensions were a staple of my SSIS years. The SCD wizard (which everyone eventually replaced with hand-written T-SQL MERGE), the hash comparisons, the IsCurrent flags. Since then I've written the same logic with Delta Lake in Azure Databricks and as Spark SQL history tables in AWS Glue. Moving to a lakehouse, I was glad to find that Delta Lake's MERGE INTO handles both Type 1 and Type 2 cleanly in Spark SQL. This post walks through both, including the "staged union" trick that does SCD2 in a single statement.

In delegated authority, the coverholder dimension is the classic SCD2 case. A coverholder's Lloyd's approval status, oversight rating, country or trading name can all change during a year of account, and when you analyse a risk or a claim you need to know what the coverholder looked like at that time. The broker dimension, on the other hand, is usually fine as Type 1.

Everything here works on Delta tables in Microsoft Fabric, Databricks, or open-source Spark with Delta Lake.

Setup

 1: CREATE TABLE IF NOT EXISTS dim_coverholder (
 2:   coverholder_sk      BIGINT GENERATED ALWAYS AS IDENTITY,
 3:   coverholder_id      STRING,      -- Lloyd's coverholder PIN
 4:   coverholder_name    STRING,
 5:   country             STRING,
 6:   approval_status     STRING,      -- APPROVED, SUSPENDED, TERMINATED
 7:   oversight_rating    STRING,      -- e.g. LOW / MEDIUM / HIGH risk
 8:   row_hash            STRING,
 9:   valid_from          TIMESTAMP,
10:   valid_to            TIMESTAMP,
11:   is_current          BOOLEAN
12: ) USING DELTA;

A note on GENERATED ALWAYS AS IDENTITY: it's supported on Databricks and in recent Delta Lake releases (3.x), but not every runtime supports it. If yours doesn't, generate surrogate keys with monotonically_increasing_id() plus the current max, or use a hash of the business key and valid_from. Identity columns also limit concurrent writes, so test that before relying on them.

The incoming batch from the coverholder register:

1: CREATE OR REPLACE TEMP VIEW stg_coverholder AS
2: SELECT coverholder_id, coverholder_name, country, approval_status, oversight_rating,
3:        sha2(concat_ws('||',
4:               coalesce(coverholder_name, '~'),
5:               coalesce(country, '~'),
6:               coalesce(approval_status, '~'),
7:               coalesce(oversight_rating, '~')), 256) AS row_hash
8: FROM   bronze.coverholder_register;

I hash the tracked attributes once in staging, so the comparison is a single column rather than a long chain of OR a.x <> b.x conditions. The coalesce matters: concat_ws silently skips nulls, so without a placeholder, ('A', NULL, 'B') and (NULL, 'A', 'B') would hash the same. It also avoids the classic trap where NULL <> 'x' evaluates to unknown and a real change gets missed.

SCD Type 1: the broker dimension

For brokers we only care about current values. A corrected broker name should simply overwrite the old one.

 1: MERGE INTO dim_broker AS t
 2: USING stg_broker      AS s
 3: ON  t.broker_code = s.broker_code
 4: WHEN MATCHED AND t.row_hash <> s.row_hash THEN UPDATE SET
 5:   t.broker_name    = s.broker_name,
 6:   t.broker_group   = s.broker_group,
 7:   t.country        = s.country,
 8:   t.row_hash       = s.row_hash
 9: WHEN NOT MATCHED THEN INSERT
10:   (broker_code, broker_name, broker_group, country, row_hash)
11:   VALUES
12:   (s.broker_code, s.broker_name, s.broker_group, s.country, s.row_hash);

The AND t.row_hash <> s.row_hash on the matched branch matters. Without it, every matched row is rewritten, which in Delta means rewriting whole Parquet files and bloating the transaction log.

SCD Type 2: the problem

For a changed coverholder, say one that's been suspended, Type 2 needs two actions: close the current row (set valid_to and is_current = false) and insert a new current row with the new status. But a single source row can only fire one MERGE branch. That's why you often see SCD2 done as an UPDATE followed by an INSERT, with a window between them where the dimension is inconsistent.

The staged-union trick

Feed MERGE each changed row twice: once with its real key (to match and close the old row) and once with a NULL merge key (which can never match, so it falls through to INSERT).

 1: MERGE INTO dim_coverholder AS t
 2: USING (
 3:   -- 1) All staging rows, keyed normally: they close changed rows / insert new coverholders
 4:   SELECT s.coverholder_id AS merge_key, s.*
 5:   FROM   stg_coverholder s
 6: 
 7:   UNION ALL
 8: 
 9:   -- 2) Changed rows again with a NULL key: they always become new current versions
10:   SELECT NULL AS merge_key, s.*
11:   FROM   stg_coverholder s
12:   JOIN   dim_coverholder t
13:     ON   t.coverholder_id = s.coverholder_id
14:    AND   t.is_current = true
15:   WHERE  t.row_hash <> s.row_hash
16: ) AS src
17: ON  t.coverholder_id = src.merge_key
18: AND t.is_current     = true
19: 
20: WHEN MATCHED AND t.row_hash <> src.row_hash THEN UPDATE SET
21:   t.valid_to   = current_timestamp(),
22:   t.is_current = false
23: 
24: WHEN NOT MATCHED THEN INSERT
25:   (coverholder_id, coverholder_name, country, approval_status, oversight_rating,
26:    row_hash, valid_from, valid_to, is_current)
27:   VALUES
28:   (src.coverholder_id, src.coverholder_name, src.country, src.approval_status,
29:    src.oversight_rating, src.row_hash, current_timestamp(), NULL, true);

Here's what happens to each type of row:

  • Newly approved coverholder: appears once with its real key, finds no match, and is inserted.
  • Unchanged coverholder: matches, but the hash condition is false, so nothing happens.
  • Changed coverholder (suspended, re-rated, renamed): the keyed copy matches and closes the old row. The NULL-keyed copy inserts the new version.

Because it's one MERGE, it's one Delta transaction. Readers see either the old state or the new one, never something in between.

Using it from the facts

When loading risk, premium or claim transactions, look up the coverholder version that was valid on the transaction date, not just the current one:

1: SELECT r.*, c.coverholder_sk
2: FROM   silver.risk_transaction r
3: JOIN   dim_coverholder c
4:   ON   c.coverholder_id = r.coverholder_id
5:  AND   r.risk_inception_date >= CAST(c.valid_from AS DATE)
6:  AND   (c.valid_to IS NULL OR r.risk_inception_date < CAST(c.valid_to AS DATE));

Now the question "how much premium did we write through coverholders while they were suspended?" has a straight answer.

One practical point: if you're loading history, set valid_from from the register's effective date rather than current_timestamp(), or every historical risk will fall before the first version of its coverholder.

Handling deletes

Delta Lake 2.3 and later support WHEN NOT MATCHED BY SOURCE, which is useful when you receive a full extract of the coverholder register and want to close coverholders that have disappeared:

1: WHEN NOT MATCHED BY SOURCE AND t.is_current = true THEN UPDATE SET
2:   t.valid_to   = current_timestamp(),
3:   t.is_current = false

Only use this with full extracts. With an incremental feed, every coverholder missing from today's batch would be closed. It's an easy mistake to make and an unpleasant one to discover.

Common errors

"Multiple source rows matched and attempted to modify the same target row." Your staging data has duplicate coverholder IDs, for example one coverholder trading under two names. Deduplicate first:

1: CREATE OR REPLACE TEMP VIEW stg_coverholder_dedup AS
2: SELECT * FROM (
3:   SELECT *, ROW_NUMBER() OVER (PARTITION BY coverholder_id ORDER BY extract_ts DESC) AS rn
4:   FROM bronze.coverholder_register
5: ) WHERE rn = 1;

Slow merges on large dimensions. MERGE has to find candidate files in the target. Help it:

  • Add the is_current = true filter to the ON clause (done above) so history rows are skipped.
  • Cluster the target on the merge key: Z-ORDER or liquid clustering on Databricks, or regular OPTIMIZE with V-Order in Fabric.
  • If the target is partitioned, add a partition predicate to ON so Spark can prune.

Checking the result

1: -- Every coverholder should have exactly one current row
2: SELECT coverholder_id, COUNT(*) AS current_rows
3: FROM   dim_coverholder
4: WHERE  is_current = true
5: GROUP  BY coverholder_id
6: HAVING COUNT(*) <> 1;
7: 
8: -- What changed in the last merge?
9: DESCRIBE HISTORY dim_coverholder LIMIT 5;

DESCRIBE HISTORY shows numTargetRowsInserted, numTargetRowsUpdated and so on in operationMetrics. Logging those to an audit table after every load is well worth it. It's the lakehouse version of the row counts we used to capture in SSIS.

The patterns from the SQL Server years still apply. The staged-union MERGE just does SCD2 more cleanly than we ever managed with a wizard.

No comments:

Post a Comment