Skip to content
Good engineers know 4 min read · Delta 3 practice problems ↓

Streaming into Delta tables

Write a stream into a Delta table exactly once, and read a Delta table as a stream.

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

Comfortable with these? Read on.

TL;DR Delta is both a streaming sink and a streaming source. As a sink, each micro-batch is one commit, and Delta records the batch id so a retried batch is not written twice. As a source, a stream reads new commits in order, which is how bronze, silver and gold tables chain together. Updates and deletes in a source table break that, unless you skip them or read the Change Data Feed.

Delta as a sink

PySpark · Append a stream to a table
(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

PySpark · Bronze to silver, incrementally
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 startingVersion or startingTimestamp to start from a point in history.
  • maxFilesPerTrigger and maxBytesPerTrigger limit 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.

OptionEffectUse when
skipChangeCommits=trueIgnores commits that rewrite filesDeletes 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 rowsDownstream must reflect updates and deletes
Restart from a new checkpointReprocess from a chosen versionOne-off rewrites like a schema fix

Upserting a stream

Step 1 · the feed

A stream of customer changes arrives: several versions of the same customer can be in one micro-batch.

Step 2 · dedup per batch

Inside foreachBatch, keep the latest change per key, because MERGE fails if two source rows match the same target row.

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

MERGE the batch into the target: update matches, insert new keys, delete on a delete flag.

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 →

MERGE with "latest wins" gives the same result if a batch is re-run, so retries are safe.

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.
  • VACUUM must 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

The stream fails. Use skipChangeCommits or the Change Data Feed.

MERGE without deduplicating the micro-batch

Fails when several source rows match one target row.

Aggressive VACUUM with lagging streams

Files a stream still needs are deleted.

Never compacting a streaming table

Thousands of small files per day.

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 questions

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

Solve: Apply a CDC Batch (MERGE Semantics) →

Keep going

Up next · lesson 16 of 26 · 4 min read
Partitioning done right
When to partition a table, how to choose the column, and how over-partitioning backfires.

Related lessons

Previous: Change Data Feed

Primary sources: Delta table streaming reads and writes · Structured Streaming guide