Streaming into Delta tables
Write a stream into a Delta table exactly once, and read a Delta table as a stream.
On this page
You will learn
- How to write a stream into a Delta table exactly once
- How to read a Delta table as a stream, and what breaks it
- How to upsert a stream with foreachBatch and MERGE
- How to keep streaming tables fast with compaction
Read first
- Structured Streaming basics · 4 min read
- The Delta transaction log · 8 min read
Comfortable with these? Read on.
Delta as a sink
(events.writeStream .format("delta") .outputMode("append") .option("checkpointLocation", "s3://chk/bronze_events/") .trigger(processingTime="1 minute") .toTable("bronze.events"))
Every micro-batch becomes one commit in the transaction logtransaction log: The ordered list of commits that defines which files make up a Delta table at each version. Learn more →, written atomically. The commit stores the query id and batch id (txn action); if the same batch is retried after a crash, Delta sees it was already committed and skips it. Together with the checkpointcheckpoint: In streaming, the folder where Spark records which input it has processed and its state, so a restarted query continues where it stopped. Learn more →, that gives exactly-once output.
Delta as a source
silver = (spark.readStream .table("bronze.events") .filter(F.col("event_type").isNotNull()) .dropDuplicatesWithinWatermark(["event_id"])) silver.writeStream.format("delta") \ .option("checkpointLocation", "s3://chk/silver_events/") \ .trigger(availableNow=True) \ .toTable("silver.events")
- The stream reads the current snapshotsnapshot: A complete, consistent version of a table at one moment. Readers always see one snapshot, never a half-finished write. Learn more → first, then each new commit in order. Use
startingVersionorstartingTimestampto start from a point in history. maxFilesPerTriggerandmaxBytesPerTriggerlimit how much each batch reads, useful for a first run over a big table.- With
trigger(availableNow=True)on a schedule, the same code is an incremental batch job: it processes only commits added since the last run.
When the source table is updated or deleted
A Delta streaming source expects append-only commits. If someone runs UPDATE, DELETE, MERGE or overwrites a partition in the source table, the stream fails, because rewritten files contain rows it has already processed.
| Option | Effect | Use when |
|---|---|---|
skipChangeCommits=true | Ignores commits that rewrite files | Deletes for retention or GDPR that downstream does not need to mirror |
| Read the Change Data Feed | .option("readChangeFeed", "true") streams inserts, updates and deletes as rows | Downstream must reflect updates and deletes |
| Restart from a new checkpoint | Reprocess from a chosen version | One-off rewrites like a schema fix |
Upserting a stream
Step 1 · the feed
Step 2 · dedup per batch
The command
w = Window.partitionBy("customer_id").orderBy(F.desc("changed_at")) latest = batch_df.withColumn("row_num", F.row_number().over(w)).filter("row_num = 1").drop("row_num")
Step 3 · MERGE
The command
(DeltaTable.forName(spark, "silver.customers").alias("t") .merge(latest.alias("s"), "t.customer_id = s.customer_id") .whenMatchedDelete(condition="s.op = 'D'") .whenMatchedUpdateAll() .whenNotMatchedInsertAll(condition="s.op != 'D'") .execute())
Step 4 · idempotentidempotent: Safe to run again: running a step twice gives the same result as running it once. Learn more →
Keeping streaming tables fast
- A one-minute trigger writes 1,440 commits a day, each with small files. Schedule
OPTIMIZE(compaction is safe while the stream runs) or enable auto compaction and optimized writes. - Delta writes a log checkpoint every 10 commits by default, so readers do not replay thousands of JSON commits.
VACUUMmust keep files at least as long as any downstream stream could lag behind, or the lagging stream fails when its files are gone.- Longer triggers or availableNow schedules reduce commits and files when minute-level freshness is not needed.
Common mistakes
UPDATE or DELETE on a table that a stream reads
MERGE without deduplicating the micro-batch
Aggressive VACUUM with lagging streams
Never compacting a streaming table
Key takeaways
- Each micro-batch is one Delta commit; batch ids make retries safe.
- Delta streaming sources expect appends; updates and deletes need skipChangeCommits or CDF.
- foreachBatch + dedup + MERGE upserts a stream idempotently.
- Compact streaming tables and keep VACUUM retention above stream lag.
Check yourself
3 questions1. How does a Delta sink avoid writing a retried micro-batch twice?
Show the answer
It records the query and batch id in the commit and skips a batch already committed. The txn action stores the batch id per query.
2. A stream reads table A. Someone runs DELETE on A. What happens by default?
Show the answer
The stream fails. Rewritten files break the append-only expectation; use skipChangeCommits or CDF.
3. Why deduplicate inside foreachBatch before MERGE?
Show the answer
MERGE fails when several source rows match one target row. A target row can only be updated by one source row per MERGE.
Practice it
Interview problems that use this: write the PySpark, run it, and get graded on hidden tests.
Keep going
Up next · lesson 16 of 26 · 4 min readPartitioning done right
When to partition a table, how to choose the column, and how over-partitioning backfires.
Related lessons
Change Data FeedData lake & lakehouse · 4 min read
MERGE INTO and upsertsSpark internals · 4 min read
Watermarks and late data
Previous: Change Data Feed
Primary sources: Delta table streaming reads and writes · Structured Streaming guide