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

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 and 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:

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, load incoming records and enforce schema:


# 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:

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:

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:


# 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:


# 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:


# 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.

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 →