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 INTOstatement or DataFrame API'sdeltaTable.merge()to upsert records atomically. - Compact and vacuum – Run
OPTIMIZEto coalesce small files andVACUUMto 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 INTOoperation viaDeltaTable.merge()to perform atomic upserts with exactly-once semantics - Configure Spark with
DeltaSparkSessionExtensionandDeltaCatalogto enable Delta functionality - Implement streaming incremental pipelines using
foreachBatchto apply merge logic on micro-batches - Maintain table performance with
OPTIMIZE(file compaction) andVACUUM(old version cleanup) - Enable time-travel debugging using
versionAsOfoptions 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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →