How to Implement Incremental Loading for SCD Tables: A Complete Spark SQL Guide
Incremental loading for Slowly Changing Dimension (SCD) tables merges existing historical records with only newly arriving or changed rows for each period, preserving full history through start_season and end_season versioning.
The DataExpert-io/data-engineer-handbook repository provides a production-ready reference implementation for incremental loading for SCD tables using Spark SQL. The pattern targets a Type-2 SCD structure for the players dataset, demonstrating how to build an initial snapshot and then efficiently merge subsequent seasons without full recomputation.
SCD Table Structure and Components
The implementation centers on four core components that work together to maintain historical accuracy while minimizing processing overhead.
| Component | Purpose | Source File |
|---|---|---|
| Raw players table | Current-season source data | players_scd_job.py |
| players_scd table | Type-2 dimension with scoring_class, is_active, start_season, end_season |
players_scd_job.py |
| Back-fill query | One-time population of full historical SCD | players_scd_table.sql |
| Incremental query | Period-by-period merge of new data with existing SCD | incremental_scd_query.sql |
The dimensional model tracks when a player's scoring_class or is_active status changes, creating new rows for each distinct period rather than overwriting previous values.
Building the Initial SCD Snapshot
Before incremental loading can begin, you need a base Type-2 SCD table. The do_player_scd_transformation function in players_scd_job.py generates this using streak detection to collapse consecutive seasons with identical attributes into single rows.
from pyspark.sql import SparkSession
query = """
WITH streak_started AS (
SELECT player_name,
current_season,
scoring_class,
LAG(scoring_class) OVER (PARTITION BY player_name ORDER BY current_season) <> scoring_class
OR LAG(scoring_class) OVER (PARTITION BY player_name ORDER BY current_season) IS NULL AS did_change
FROM players
),
streak_identified AS (
SELECT player_name,
scoring_class,
current_season,
SUM(CASE WHEN did_change THEN 1 ELSE 0 END)
OVER (PARTITION BY player_name ORDER BY current_season) AS streak_identifier
FROM streak_started
),
aggregated AS (
SELECT player_name,
scoring_class,
streak_identifier,
MIN(current_season) AS start_date,
MAX(current_season) AS end_date
FROM streak_identified
GROUP BY 1,2,3
)
SELECT player_name, scoring_class, start_date, end_date
FROM aggregated
"""
def do_player_scd_transformation(spark, dataframe):
dataframe.createOrReplaceTempView("players")
return spark.sql(query)
if __name__ == "__main__":
spark = SparkSession.builder.master("local").appName("players_scd").getOrCreate()
output_df = do_player_scd_transformation(spark, spark.table("players"))
output_df.write.mode("overwrite").insertInto("players_scd")
The query logic proceeds through three CTE stages:
- streak_started: Identifies season boundaries where
scoring_classchanges usingLAGwindow function - streak_identified: Assigns a cumulative
streak_identifierto group consecutive identical values - aggregated: Compresses each streak into a single row with
start_dateandend_date
The Incremental Loading Query
The core incremental loading for SCD tables logic lives in incremental_scd_query.sql. This query handles four distinct record categories in a single execution:
WITH last_season_scd AS (
SELECT * FROM players_scd
WHERE current_season = 2021
AND end_season = 2021
),
historical_scd AS (
SELECT player_name,
scoring_class,
is_active,
start_season,
end_season
FROM players_scd
WHERE current_season = 2021
AND end_season < 2021
),
this_season_data AS (
SELECT * FROM players
WHERE current_season = 2022
),
unchanged_records AS (
SELECT ts.player_name,
ts.scoring_class,
ts.is_active,
ls.start_season,
ts.current_season AS end_season
FROM this_season_data ts
JOIN last_season_scd ls ON ls.player_name = ts.player_name
WHERE ts.scoring_class = ls.scoring_class
AND ts.is_active = ls.is_active
),
changed_records AS (
SELECT ts.player_name,
UNNEST(ARRAY[
ROW(ls.scoring_class, ls.is_active, ls.start_season, ls.end_season)::scd_type,
ROW(ts.scoring_class, ts.is_active, ts.current_season, ts.current_season)::scd_type
]) AS records
FROM this_season_data ts
LEFT JOIN last_season_scd ls ON ls.player_name = ts.player_name
WHERE (ts.scoring_class <> ls.scoring_class
OR ts.is_active <> ls.is_active)
),
unnested_changed_records AS (
SELECT player_name,
(records::scd_type).scoring_class,
(records::scd_type).is_active,
(records::scd_type).start_season,
(records::scd_type).end_season
FROM changed_records
),
new_records AS (
SELECT ts.player_name,
ts.scoring_class,
ts.is_active,
ts.current_season AS start_season,
ts.current_season AS end_season
FROM this_season_data ts
LEFT JOIN last_season_scd ls ON ts.player_name = ls.player_name
WHERE ls.player_name IS NULL
)
SELECT *, 2022 AS current_season FROM (
SELECT * FROM historical_scd
UNION ALL
SELECT * FROM unchanged_records
UNION ALL
SELECT * FROM unnested_changed_records
UNION ALL
SELECT * FROM new_records
) a;
Record Category Handling
| Category | Detection Logic | Action Taken |
|---|---|---|
| Unchanged | scoring_class and is_active match previous season |
Extend end_season to current season |
| Changed | Either attribute differs from previous season | Close old period, open new period (2 rows) |
| New | player_name not found in last_season_scd |
Insert single-row period starting current season |
| Historical | Previously closed records (end_season < current_season) |
Pass through unchanged |
The changed_records CTE uses UNNEST with an ARRAY of ROW constructors to generate both the closed historical record and the new open record in one operation, then explodes them in unnested_changed_records.
Key Implementation Files
The complete pattern spans four files in the repository:
players_scd_job.py: Spark job computing initial Type-2 SCD snapshots with streak-based aggregationplayers_scd_table.sql: DDL defining theplayers_scdtable structure, partitioning, and constraintsincremental_scd_query.sql: Full incremental loading for SCD tables query with CTE-based change detectiontest_player_scd.py: PyTest validation suite ensuring transformation correctness
Adapting This Pattern
To implement incremental loading for SCD tables in your own pipeline:
- Replace
player_namewith your natural business key - Modify
scoring_classandis_activeto your slowly changing attributes - Adjust
current_seasonto your grain (date, month, quarter) - Define the
scd_typecomposite type in your database for theUNNESToperation - Parameterize the season/year values instead of hardcoding
Summary
- Streak detection collapses consecutive identical records to minimize SCD row count
- Four-way CTE structure cleanly separates unchanged, changed, new, and historical records
UNNESTwith composite types efficiently generates paired close/open rows for Type-2 changes- Partition filter on
current_seasonensures the incremental query touches only relevant data - Complete source code resides in DataExpert-io/data-engineer-handbook with executable tests
Frequently Asked Questions
What is the difference between back-filling and incremental loading for SCD tables?
Back-filling processes the entire historical dataset once to establish the complete SCD history, while incremental loading processes only new or changed records for each subsequent period. In the repository's pattern, players_scd_table.sql handles the one-time back-fill, and incremental_scd_query.sql runs repeatedly for each new season.
Why use streak detection instead of season-by-season records?
Streak detection reduces storage and improves query performance by collapsing consecutive seasons with identical attributes into single rows. A player with the same scoring_class for 10 consecutive seasons generates one row instead of ten, compressing the dimension table significantly for slowly changing attributes.
How does the query handle multiple attribute changes simultaneously?
The changed_records CTE detects any difference in either scoring_class or is_active using the WHERE clause condition, then the UNNEST operation generates both the closed historical record and the new active record regardless of which specific attribute changed. This ensures proper Type-2 versioning even when multiple attributes change together.
Can this pattern work with streaming data instead of batch seasons?
The repository's implementation uses batch processing with current_season as a partitioning key. For streaming, you would replace the season-based filters with micro-batch windows and likely use MERGE INTO syntax (Databricks Delta Lake) or equivalent upsert operations instead of the UNION ALL pattern, though the core change-detection logic remains applicable.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →