You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

求助:如何在Spark SQL中实现SCD Type2+Type1,在Azure Databricks捕获新旧记录

Hey there! Let's walk through how to implement historical record capture and data change tracking in Azure Databricks Notebooks, plus a hybrid SCD Type 1 + Type 2 solution using Spark SQL. I'll use your sample data as a reference.

Understanding Your Sample Data

First, let's clean up your sample data into a readable table—this represents the change history for employee 1550:

EMPNOROW_EFF_DTSTATE_CODECITY_NAMEROW_EXPIRY_DT
155019-10-2020TAFORTWORTH20-10-2020
155020-10-2020TXRONOLE21-01-2020
155021-01-2020TXGRAPEVINNULL

The NULL in ROW_EXPIRY_DT indicates this is the current active record for the employee.

What's the Hybrid SCD Approach?

  • SCD Type 2: Tracks full history of changes for fields where you need to retain past values (like STATE_CODE and CITY_NAME here). We mark old records with an expiry date and insert new active records.
  • SCD Type 1: Overwrites old values with new ones for fields where you don't need historical context (e.g., employee status, email address)—this keeps the latest value only.

For your scenario, we'll use Type 2 for location-related fields and Type 1 for any non-historical fields you might have (like EMP_STATUS in the examples below).

Step-by-Step Implementation with Spark SQL

We'll use Delta Lake (recommended for Azure Databricks) because it supports ACID transactions, MERGE INTO operations, and time travel for historical auditing.

1. Set Up Source and Dimension Tables

First, create your source data view and dimension table:

-- Create a temp view for your source data (replace with your actual data source)
CREATE OR REPLACE TEMP VIEW employee_source AS
SELECT 1550 AS EMPNO, '19-10-2020' AS ROW_EFF_DT, 'TA' AS STATE_CODE, 'FORTWORTH' AS CITY_NAME, '20-10-2020' AS ROW_EXPIRY_DT
UNION ALL
SELECT 1550 AS EMPNO, '20-10-2020' AS ROW_EFF_DT, 'TX' AS STATE_CODE, 'RONOLE' AS CITY_NAME, '21-01-2020' AS ROW_EXPIRY_DT
UNION ALL
SELECT 1550 AS EMPNO, '21-01-2020' AS ROW_EFF_DT, 'TX' AS STATE_CODE, 'GRAPEVIN' AS CITY_NAME, NULL AS ROW_EXPIRY_DT;

-- Create Delta table for employee dimension (stores SCD data)
CREATE OR REPLACE TABLE employee_dim (
    EMPNO INT,
    STATE_CODE STRING,
    CITY_NAME STRING,
    ROW_EFF_DT DATE,
    ROW_EXPIRY_DT DATE,
    IS_CURRENT BOOLEAN,
    EMP_STATUS STRING -- Type 1 field: no history needed, overwrite on change
)
USING DELTA;

2. Clean and Prepare Source Data

Convert date strings to proper DATE types to avoid formatting issues:

CREATE OR REPLACE TEMP VIEW employee_source_clean AS
SELECT
    EMPNO,
    TO_DATE(ROW_EFF_DT, 'dd-MM-yyyy') AS ROW_EFF_DT,
    STATE_CODE,
    CITY_NAME,
    CASE WHEN ROW_EXPIRY_DT IS NULL THEN NULL ELSE TO_DATE(ROW_EXPIRY_DT, 'dd-MM-yyyy') END AS ROW_EXPIRY_DT
FROM employee_source;

3. Identify Changes

Detect which records need Type 2 updates (location changes) and Type 1 updates (status changes):

-- Get current active records from the dimension table
CREATE OR REPLACE TEMP VIEW current_active_employees AS
SELECT * FROM employee_dim WHERE IS_CURRENT = true;

-- Find records where Type 2 fields (location) have changed
CREATE OR REPLACE TEMP VIEW type2_changed_records AS
SELECT
    s.EMPNO,
    s.STATE_CODE,
    s.CITY_NAME,
    s.ROW_EFF_DT,
    s.ROW_EXPIRY_DT
FROM employee_source_clean s
JOIN current_active_employees c
    ON s.EMPNO = c.EMPNO
WHERE s.STATE_CODE != c.STATE_CODE OR s.CITY_NAME != c.CITY_NAME;

-- Prepare new records to insert (mark as current)
CREATE OR REPLACE TEMP VIEW new_records_to_insert AS
SELECT
    EMPNO,
    STATE_CODE,
    CITY_NAME,
    ROW_EFF_DT,
    ROW_EXPIRY_DT,
    TRUE AS IS_CURRENT,
    'ACTIVE' AS EMP_STATUS -- Default Type 1 value
FROM employee_source_clean
WHERE EMPNO NOT IN (SELECT EMPNO FROM current_active_employees)
UNION ALL
SELECT
    s.EMPNO,
    s.STATE_CODE,
    s.CITY_NAME,
    s.ROW_EFF_DT,
    s.ROW_EXPIRY_DT,
    TRUE AS IS_CURRENT,
    c.EMP_STATUS -- Retain existing Type 1 value unless updating
FROM type2_changed_records s
JOIN current_active_employees c ON s.EMPNO = c.EMPNO;

4. Execute MERGE Operation

This single command handles:

  • Expiring old Type 2 records
  • Inserting new active records
  • Updating Type 1 fields
MERGE INTO employee_dim target
USING (
    -- Combine records to expire, new records to insert, and Type 1 updates
    SELECT
        c.EMPNO,
        c.STATE_CODE,
        c.CITY_NAME,
        c.ROW_EFF_DT,
        s.ROW_EFF_DT AS NEW_ROW_EXPIRY_DT,
        FALSE AS IS_CURRENT,
        c.EMP_STATUS
    FROM current_active_employees c
    JOIN type2_changed_records s ON c.EMPNO = s.EMPNO
    UNION ALL
    SELECT * FROM new_records_to_insert
    UNION ALL
    -- Type 1 update: Overwrite EMP_STATUS if it changed
    SELECT
        s.EMPNO,
        c.STATE_CODE,
        c.CITY_NAME,
        c.ROW_EFF_DT,
        c.ROW_EXPIRY_DT,
        TRUE AS IS_CURRENT,
        s.EMP_STATUS -- New Type 1 value
    FROM employee_source_clean s
    JOIN current_active_employees c ON s.EMPNO = c.EMPNO
    WHERE s.EMP_STATUS != c.EMP_STATUS
) source
ON target.EMPNO = source.EMPNO AND target.IS_CURRENT = true
WHEN MATCHED AND source.IS_CURRENT = false THEN
    UPDATE SET target.ROW_EXPIRY_DT = source.NEW_ROW_EXPIRY_DT, target.IS_CURRENT = false
WHEN MATCHED AND source.IS_CURRENT = true THEN
    UPDATE SET target.EMP_STATUS = source.EMP_STATUS
WHEN NOT MATCHED THEN
    INSERT (EMPNO, STATE_CODE, CITY_NAME, ROW_EFF_DT, ROW_EXPIRY_DT, IS_CURRENT, EMP_STATUS)
    VALUES (source.EMPNO, source.STATE_CODE, source.CITY_NAME, source.ROW_EFF_DT, source.ROW_EXPIRY_DT, source.IS_CURRENT, source.EMP_STATUS);

5. Verify the Result

Run this query to check your dimension table:

SELECT * FROM employee_dim ORDER BY EMPNO, ROW_EFF_DT;

You should see all historical records with correct expiry dates, and the latest record marked as IS_CURRENT = true.

Key Technical Recommendations

  • Always Use Delta Lake: It simplifies SCD implementation with MERGE INTO, ensures data consistency, and lets you use time travel (TIMESTAMP AS OF) to audit past data.
  • Clear Field Classification: Explicitly define which fields are Type 1 vs Type 2 to avoid mistakes. For example, use Type 1 for fields like EMP_STATUS or PHONE_NUMBER where history isn't needed.
  • Incremental Processing: For large datasets, process only new/changed records instead of full loads. Use Databricks Auto Loader to incrementally ingest data from sources like Blob Storage or ADLS.
  • Date Validation: Ensure all date fields are stored as DATE types, not strings, to prevent conversion errors and improve query performance.
  • Index Optimization: Add indexes on EMPNO and IS_CURRENT to speed up joins and merge operations for large dimension tables.

内容的提问来源于stack exchange,提问作者shoyab

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.09 15:42:46