# How to Implement Incremental Loading for SCD Tables: A Complete Spark SQL Guide

> Learn how to implement incremental loading for SCD tables using Spark SQL. Merge new and changed rows with existing history using start and end dates for complete versioning.

- Repository: [DataExpert.io/data-engineer-handbook](https://github.com/DataExpert-io/data-engineer-handbook)
- Tags: how-to-guide
- Published: 2026-08-06

---

**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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/players_scd_job.py) |
| **players_scd table** | Type-2 dimension with `scoring_class`, `is_active`, `start_season`, `end_season` | [`players_scd_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/players_scd_job.py) |
| **Back-fill query** | One-time population of full historical SCD | [`players_scd_table.sql`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/players_scd_table.sql) |
| **Incremental query** | Period-by-period merge of new data with existing SCD | [`incremental_scd_query.sql`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/players_scd_job.py) generates this using **streak detection** to collapse consecutive seasons with identical attributes into single rows.

```python
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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/incremental_scd_query.sql). This query handles four distinct record categories in a single execution:

```sql
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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/players_scd_job.py)**: Spark job computing initial Type-2 SCD snapshots with streak-based aggregation
- **[`players_scd_table.sql`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/players_scd_table.sql)**: DDL defining the `players_scd` table structure, partitioning, and constraints
- **[`incremental_scd_query.sql`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/incremental_scd_query.sql)**: Full **incremental loading for SCD tables** query with CTE-based change detection
- **[`test_player_scd.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/test_player_scd.py)**: PyTest validation suite ensuring transformation correctness

## 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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/players_scd_table.sql) handles the one-time back-fill, and [`incremental_scd_query.sql`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/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.