求助:如何在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:
| EMPNO | ROW_EFF_DT | STATE_CODE | CITY_NAME | ROW_EXPIRY_DT |
|---|---|---|---|---|
| 1550 | 19-10-2020 | TA | FORTWORTH | 20-10-2020 |
| 1550 | 20-10-2020 | TX | RONOLE | 21-01-2020 |
| 1550 | 21-01-2020 | TX | GRAPEVIN | NULL |
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_CODEandCITY_NAMEhere). 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_STATUSorPHONE_NUMBERwhere 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
DATEtypes, not strings, to prevent conversion errors and improve query performance. - Index Optimization: Add indexes on
EMPNOandIS_CURRENTto speed up joins and merge operations for large dimension tables.
内容的提问来源于stack exchange,提问作者shoyab

