# How to Implement Incremental Processing with Delta Lake: A Complete Guide

> Learn to implement incremental processing with Delta Lake using MERGE INTO for efficient data upserts. Achieve exactly-once semantics and time-travel for reliable data pipelines.

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

---

**Implement incremental processing with Delta Lake by using ACID transactions and the MERGE INTO operation to upsert data, enabling exactly-once semantics and time-travel capabilities for reliable data pipelines.**

Delta Lake adds ACID transactions, versioning, and schema enforcement on top of Apache Spark SQL, making it ideal for building reliable incremental data pipelines. According to the DataExpert-io/data-engineer-handbook repository, you can extend standard Spark batch jobs—like those in [`team_vertex_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/team_vertex_job.py) and [`players_scd_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/players_scd_job.py)—to support incremental processing with exactly-once guarantees. This guide demonstrates how to implement incremental processing with Delta Lake using production-ready patterns derived from the handbook's Spark fundamentals.

## Core Architecture for Incremental Processing

The typical architecture for incremental processing with Delta Lake follows three atomic steps:

- **Ingest new data** – Load raw or change-data-capture (CDC) records into a Spark DataFrame.
- **Merge into the target Delta table** – Use the `MERGE INTO` statement or DataFrame API's `deltaTable.merge()` to upsert records atomically.
- **Compact and vacuum** – Run `OPTIMIZE` to coalesce small files and `VACUUM` to remove old versions.

Because Delta Lake tracks table versions, you can also **time-travel** to prior snapshots for debugging or re-processing.

## Configuring Spark for Delta Lake

Before implementing incremental logic, configure your Spark session with Delta Lake extensions. This pattern mirrors the initialization found in [`team_vertex_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/team_vertex_job.py):

```python
from pyspark.sql import SparkSession

spark = (
    SparkSession.builder
        .appName("IncrementalDeltaPipeline")
        .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")
        .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog")
        .getOrCreate()
)

```

## Implementing Incremental Upserts

### Loading New Data

Following the data ingestion patterns from [`monthly_user_site_hits_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/monthly_user_site_hits_job.py), load incoming records and enforce schema:

```python

# Load new records (e.g., from a JSON source)

new_df = spark.read.json("s3://my-bucket/incoming/2024-08-09/*.json")

# Schema enforcement via selectExpr

new_df = new_df.selectExpr(
    "cast(id as string) as id",
    "cast(event_ts as timestamp) as event_ts",
    "payload"
)

```

### Defining the Target Delta Table

Create or reference your Delta table. For initial loads, use append mode; subsequent runs use merge:

```python
delta_path = "s3://my-bucket/delta/events"

# Initial load only - creates the table

new_df.write.format("delta").mode("append").save(delta_path)

```

### Executing Atomic Merges

The core of incremental processing uses **DeltaTable's merge API**. This approach extends the SCD (Slowly Changing Dimension) logic demonstrated in [`players_scd_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/players_scd_job.py):

```python
from delta.tables import DeltaTable

delta_tbl = DeltaTable.forPath(spark, delta_path)

(
    delta_tbl.alias("t")
    .merge(
        new_df.alias("s"),
        "t.id = s.id"
    )
    .whenMatchedUpdateAll()
    .whenNotMatchedInsertAll()
    .execute()
)

```

This operation atomically updates existing records where IDs match and inserts new records where they don't.

## Streaming Incremental Processing

For near real-time workflows, combine Delta Lake with **Structured Streaming** using `foreachBatch`. This applies the same merge logic on micro-batches:

```python

# Streaming source (Kafka, Pub/Sub, etc.)

new_stream = spark.readStream.format("json").schema(new_df.schema).load("s3://my-bucket/incoming/stream/")

(
    new_stream.writeStream
    .foreachBatch(lambda batch_df, batch_id: 
        delta_tbl.alias("t")
            .merge(
                batch_df.alias("s"),
                "t.id = s.id"
            )
            .whenMatchedUpdateAll()
            .whenNotMatchedInsertAll()
            .execute()
    )
    .outputMode("update")
    .start()
)

```

## Maintenance and Time Travel

### Optimizing Storage

Periodically run maintenance operations to keep storage efficient:

```python

# Coalesce small files into larger ones

spark.sql(f"OPTIMIZE delta.`{delta_path}`")

# Remove old versions (retention default is 7 days)

spark.sql(f"VACUUM delta.`{delta_path}` RETAIN 168 HOURS")

```

### Querying Historical Versions

Leverage Delta's time-travel for debugging or re-processing specific versions:

```python

# Read table as it existed 2 versions ago

historical_df = spark.read.format("delta").option("versionAsOf", 2).load(delta_path)
historical_df.show()

```

## Summary

- **Delta Lake** provides ACID transactions and versioning on top of Spark SQL for reliable incremental processing
- Use the **`MERGE INTO`** operation via `DeltaTable.merge()` to perform atomic upserts with exactly-once semantics
- Configure Spark with `DeltaSparkSessionExtension` and `DeltaCatalog` to enable Delta functionality
- Implement streaming incremental pipelines using **`foreachBatch`** to apply merge logic on micro-batches
- Maintain table performance with **`OPTIMIZE`** (file compaction) and **`VACUUM`** (old version cleanup)
- Enable time-travel debugging using **`versionAsOf`** options when reading Delta tables

## Frequently Asked Questions

### What is the difference between Delta Lake incremental processing and standard Spark append mode?

Standard Spark append mode only adds new records, risking duplicates if jobs retry. Delta Lake's merge operation provides atomic upserts, updating existing records and inserting new ones within a transaction, ensuring exactly-once processing even during failures.

### How does Delta Lake handle schema evolution during incremental processing?

Delta Lake enforces schema by default but supports schema evolution through the `mergeSchema` option when writing, or by enabling `spark.databricks.delta.schema.autoMerge.enabled`. During merges, you can explicitly cast fields to handle incoming data with different types, as shown in the repository's data ingestion patterns.

### Can I use Delta Lake incremental processing with streaming sources like Kafka?

Yes. Use Structured Streaming's `readStream` to consume from Kafka, then apply the merge logic within `foreachBatch`. This converts micro-batches into atomic Delta Lake transactions, providing exactly-once semantics across streaming and batch workflows.

### How often should I run OPTIMIZE and VACUUM on Delta tables?

Run `OPTIMIZE` daily or when small files accumulate (typically after many incremental writes), and `VACUUM` based on your retention requirements (default 7 days). Never vacuum immediately after optimizing if you need time-travel capability for recent history.