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_class changes using LAG window function
  • streak_identified: Assigns a cumulative streak_identifier to group consecutive identical values
  • aggregated: Compresses each streak into a single row with start_date and end_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:

Adapting This Pattern

To implement incremental loading for SCD tables in your own pipeline:

  1. Replace player_name with your natural business key
  2. Modify scoring_class and is_active to your slowly changing attributes
  3. Adjust current_season to your grain (date, month, quarter)
  4. Define the scd_type composite type in your database for the UNNEST operation
  5. 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
  • UNNEST with composite types efficiently generates paired close/open rows for Type-2 changes
  • Partition filter on current_season ensures 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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →