foreachBatch is the escape hatch of Spark Structured Streaming. It hands you every micro-batch as a plain batch DataFrame together with a batch identifier, and lets you do anything a batch job can do with it: MERGE into a Delta table, upsert into PostgreSQL, write to three places at once, or call a library that has no streaming sink. That freedom is why it shows up in most production streaming jobs, and also why it is the source of most duplicate-row incidents.

This article explains the contract, what the batchId promises and why the default is at-least-once, then builds a Kafka pipeline that upserts orders into Delta and PostgreSQL exactly once in effect, and ends with sizing, failure modes and a checklist.

Advertisement

The contract in one paragraph

You register a function with df.writeStream.foreachBatch(fn). For each micro-batch the engine plans an incremental query over the new input, materialises it as a DataFrame that is not a streaming DataFrame, and calls fn(batch_df, batch_id) on the driver. Its actions run as normal Spark jobs. When the function returns without an exception, the engine writes the batch id to the commit log and moves on. The API has existed since Spark 2.4 in Scala, Java and Python. It works only with micro-batch execution: continuous processing mode does not support it, and the documentation points you to foreach instead.

Three consequences follow. First, the batch DataFrame supports every batch operation, including ones streaming DataFrames reject, such as MERGE, multiple aggregations, limit or writing through a JDBC driver. Second, every action you call inside the function recomputes the batch's lineage unless you cache it, and for a stateful query that means reading state again. Third, Spark has no idea what your function wrote. It cannot roll back an external write, so the guarantee Spark documents is at-least-once: after a failure the same batch can be written again.

Where the function sits in the micro-batch lifecycle

A micro-batch has a fixed order of durable steps, and the replay behaviour comes straight from that order. The engine decides the offset range for batch N and writes it to checkpoint/offsets/N before running anything. It then runs the incremental plan, updating the state store if the query is stateful, and calls your function. Only after the function returns does it write checkpoint/commits/N.

One micro-batch, from planned offsets to committed batchIdSourceKafka, files, DeltaOffset logcheckpoint/offsets/NPlan batch Nincremental queryState storestateful operatorsforeachBatch(df, N)your function, driver sidebatch DataFrameDelta MERGEtxnAppId + txnVersionJDBC upsertON CONFLICTBatch-id ledgerskip if seenMetrics / alertsside outputsCommit logcheckpoint/commits/Nfunction returnedA crash between the function's writes and commits/N replays batch N with the same offsets and the same batchId
The offset log is written before the function runs and the commit log after it returns. Everything your function does happens in the window between those two writes.

On restart the engine compares the two logs. If the latest offset entry has no matching commit, batch N is re-run with exactly the offsets recorded in the offset log and with the same batchId. For a replayable source such as Kafka or a file source, a replayed batch reads the same input and is labelled with the same number. If your transformation is deterministic, the replayed batch DataFrame contains the same rows.

The weak spot is the gap. If the function wrote to two sinks and the driver died after the first, or the function finished but the process was killed before the commit file landed, the external systems already hold batch N. The replay writes it again. Every correct use of foreachBatch is a strategy for making that second write harmless.

Advertisement

A first pipeline: Kafka to Delta with built-in idempotency

The simplest correct pattern uses a sink that understands batch identifiers. Delta Lake's DataFrame writer accepts two options, txnAppId and txnVersion. Delta records the highest version it has committed for each application id in the transaction log, and silently skips a write whose version is not greater than that. Pass a stable query name as the application id and the batchId as the version, and a replayed batch becomes a no-op.

from pyspark.sql import SparkSession, DataFrame
from pyspark.sql import functions as F

spark = SparkSession.builder.appName("orders-upsert").getOrCreate()

orders = (spark.readStream.format("kafka")
          .option("kafka.bootstrap.servers", "kafka:9092")
          .option("subscribe", "orders")
          .option("maxOffsetsPerTrigger", 200000)
          .load()
          .select(F.from_json(F.col("value").cast("string"), ORDER_SCHEMA).alias("o"))
          .select("o.*"))

def upsert_orders(batch_df: DataFrame, batch_id: int) -> None:
    # Runs on the driver once per micro-batch; batch_df is an ordinary batch DataFrame.
    print(f"batch {batch_id}: planning write")
    (batch_df.write.format("delta").mode("append")
        .option("txnAppId", "orders-upsert-v1")   # stable per query, never per run
        .option("txnVersion", batch_id)           # Delta skips a version it has seen
        .save("/lake/silver/orders_raw"))

query = (orders.writeStream
         .foreachBatch(upsert_orders)
         .option("checkpointLocation", "/chk/orders-upsert-v1")
         .trigger(processingTime="30 seconds")
         .start())

Two details matter. The application id must stay the same across restarts of the same logical query and must differ between queries writing the same table, otherwise one query's progress hides the other's writes. And if you ever start the query with a new checkpoint, batch ids start again from zero, so you must also change the application id, which is why the example embeds a version suffix in both the id and the checkpoint path.

Three ways to make the second write harmless

Not every sink has a transaction log. Pick one of these patterns per sink, and be explicit in code review about which one each write uses.

PatternHow it worksGood forWatch out for
Sink-level batch versioningSink remembers the last batch version per writer and skips repeatsDelta via txnAppId/txnVersionMust reset the app id with a new checkpoint
Keyed upsertMERGE or INSERT ... ON CONFLICT by business key; a replay rewrites the same valuesCurrent-state tables, dimension tablesDuplicate keys inside one batch; needs an ordering column for late updates
Batch ledgerRecord (query, batchId) in the target in the same transaction as the data; skip if presentRelational databases, any transactional storeLedger and data must commit atomically, or it proves nothing

Append-only sinks without any of these properties, such as a plain Parquet directory, a message queue or an email API, cannot be made exactly-once from inside foreachBatch. Say so in the design document rather than claiming exactly-once.

Worked example: upserting orders into Delta and PostgreSQL

Suppose an orders topic carries the full current state of an order on every change. We want a Delta table holding the latest state per order for analytics, and the same latest state in PostgreSQL for a service. Each Kafka message carries order_id and updated_at.

Within one micro-batch the same order can appear several times. A Delta MERGE raises an error when more than one source row matches the same target row, so the function first keeps the newest row per key with a window. It caches that result because two sinks will read it, then runs the MERGE with an updated_at guard so an older event that arrives late cannot overwrite a newer one.

from delta.tables import DeltaTable
from pyspark.sql import Window

def merge_latest(batch_df, batch_id):
    s = batch_df.sparkSession                 # use the batch's session, not a captured global
    # 1. Collapse duplicates inside the batch: MERGE fails if two source rows hit one target row.
    w = Window.partitionBy("order_id").orderBy(F.col("updated_at").desc())
    latest = (batch_df.withColumn("rn", F.row_number().over(w))
                      .filter("rn = 1").drop("rn"))
    latest.persist()
    try:
        target = DeltaTable.forName(s, "silver.orders")
        (target.alias("t")
            .merge(latest.alias("s"), "t.order_id = s.order_id")
            .whenMatchedUpdateAll(condition="s.updated_at > t.updated_at")
            .whenNotMatchedInsertAll()
            .execute())
        # 2. Second sink from the same cached rows: a JDBC staging table, keyed by batch.
        clear_staging(batch_id)               # DELETE staging rows left by an earlier attempt
        (latest.withColumn("batch_id", F.lit(batch_id))
               .write.format("jdbc")
               .option("url", JDBC_URL).option("dbtable", "staging.orders_batch")
               .option("user", USER).option("password", PASSWORD)
               .mode("append").save())
        upsert_from_staging(batch_id)         # INSERT ... ON CONFLICT, then DELETE staging rows
    finally:
        latest.unpersist()

The MERGE is idempotent by construction: replaying the batch rewrites the same rows with the same values, and the timestamp guard makes it order-safe across batches. The PostgreSQL side uses a staging table plus a ledger, so the visible table changes exactly once per batch even though the staging insert can be repeated:

-- A ledger for sinks with no idempotency option (PostgreSQL here).
CREATE TABLE stream_batches (
  query_name text   NOT NULL,
  batch_id   bigint NOT NULL,
  rows       bigint NOT NULL,
  done_at    timestamptz NOT NULL DEFAULT now(),
  PRIMARY KEY (query_name, batch_id)
);

-- upsert_from_staging(batch_id) runs this in ONE database transaction:
BEGIN;
INSERT INTO stream_batches (query_name, batch_id, rows)
  SELECT 'orders-upsert-v1', :batch_id, count(*) FROM staging.orders_batch WHERE batch_id = :batch_id
  ON CONFLICT DO NOTHING;
-- if zero rows were inserted above, the batch already landed: DELETE its staging rows, COMMIT, return
INSERT INTO orders AS o (order_id, status, amount, updated_at)
  SELECT order_id, status, amount, updated_at FROM staging.orders_batch WHERE batch_id = :batch_id
  ON CONFLICT (order_id) DO UPDATE SET status = EXCLUDED.status, amount = EXCLUDED.amount,
      updated_at = EXCLUDED.updated_at
  WHERE o.updated_at < EXCLUDED.updated_at;
DELETE FROM staging.orders_batch WHERE batch_id = :batch_id;
COMMIT;

Walk through a crash. Batch 812 merges into Delta, appends 4,100 rows to staging and the driver is killed during the PostgreSQL transaction. On restart, batch 812 replays with the same offsets. The Delta MERGE rewrites identical values. The function first deletes the 4,100 staging rows the dead attempt left behind, appends the same 4,100 again, and the transaction runs for the first time, so the orders table changes once. Without that cleanup the staging table would hold each order twice, and PostgreSQL rejects an ON CONFLICT DO UPDATE that touches the same row twice in one statement. If the kill had come after the PostgreSQL commit, the ledger insert would find its key already present, so the procedure deletes that batch's staging rows and commits without touching orders.

Fan-out, caching and the cost of every action

Each action inside the function, whether write, count or collect, runs the batch plan again unless the DataFrame is cached. For a stateless Kafka read that means fetching the offsets twice; for a stateful aggregation it means loading and evaluating state twice, which the Spark documentation calls out as a performance problem. The pattern is persist() at the top, all writes, then unpersist() in a finally block so a failed batch does not leak cached blocks across retries.

Writes inside the function run one after another, because each action blocks the driver thread. Keep them sequential and put the write that is hardest to undo last.

After a shuffle a micro-batch often has 200 partitions, so an append writes 200 small files per trigger; coalesce first.

Sizing batches and triggers

The function's run time is part of the batch duration. In the query progress it is reported as durationMs.addBatch, next to triggerExecution for the whole batch. If addBatch regularly approaches the trigger interval, the query falls behind its source and every following batch is larger, which makes the next one slower again.

Control batch size at the source, not in the function: maxOffsetsPerTrigger for Kafka, maxFilesPerTrigger for file sources. A larger batch amortises fixed costs such as a MERGE's scan of the target, so for MERGE-heavy sinks you usually want fewer, larger batches and a longer trigger interval. For latency-sensitive sinks you want smaller batches and must make the per-batch fixed cost small, for example by partition-pruning the MERGE condition on a date column.

p = query.lastProgress            # dict for the most recent micro-batch
if p:
    print(p["batchId"], p["numInputRows"],
          p["durationMs"].get("addBatch"),        # time spent inside foreachBatch
          p["durationMs"].get("triggerExecution"),
          p["inputRowsPerSecond"], p["processedRowsPerSecond"])

# Alert when the function is the bottleneck: addBatch close to the trigger interval
# for several batches in a row means the query is falling behind its source.

Failure modes

SymptomCauseFix
Duplicate rows after a restartAppend to a sink with no idempotency; batch replayedUse txnVersion, a keyed upsert or a batch ledger
Rows silently missing after redeploySame txnAppId reused with a fresh checkpoint, so new batch ids look old to DeltaChange the app id whenever the checkpoint changes
MERGE fails: multiple source rows matchedSame key twice in one micro-batchDeduplicate inside the batch with a window first
Batch takes twice as long as expectedSeveral actions recompute the plan and reload statepersist once, unpersist in finally
Exception swallowed, data losttry/except around the write logs and returns, so the batch is committedLet exceptions propagate; Spark then fails the query and replays the batch
Function works in a notebook, fails in productionCode uses a captured global session or driver-local stateUse batch_df.sparkSession and keep state in durable stores
Output differs on replayNon-deterministic logic such as current_timestamp or random samplingDerive timestamps from the data or from batchId

A function that catches a write failure and returns normally tells Spark the batch succeeded. The commit log advances, and that batch's data never reaches the sink. Catch only to add context, then re-raise.

foreachBatch, foreach or a native sink

Prefer a native streaming sink when one does what you need. Choose foreachBatch when you need batch-only operations such as MERGE, when you write to several sinks from one read, or when the target only has a batch connector. Choose foreach with open, process and close when you need continuous mode or per-row side effects whose connection handling you want to control per partition; it is harder to make idempotent because you see rows, not batches.

The cost of foreachBatch is that correctness moves into your code. Spark guarantees each batch id is processed until it is committed; you guarantee that processing it twice changes nothing. Related reading: Kafka with Structured Streaming for offsets and delivery semantics, streaming deduplication for removing duplicates that come from producers, the state store for what a replay reloads, Delta Lake for the transaction log behind txnVersion and cache and persist for storage levels.

What to do next

  1. List every foreachBatch function in your jobs and write next to each write which idempotency pattern it uses; any write without one is at-least-once.
  2. For Delta appends, add txnAppId and txnVersion, and put the checkpoint version into the app id so a fresh checkpoint can never collide with old versions.
  3. Deduplicate by key inside each batch before any MERGE, with an ordering column guard against late updates.
  4. Wrap multi-sink functions in persist and a finally unpersist, and order writes so the hardest to undo runs last.
  5. Remove any try/except that logs and returns; re-raise after adding context.
  6. Alert when addBatch exceeds about 70 percent of the trigger interval for several consecutive batches, and tune maxOffsetsPerTrigger.
  7. Kill the driver in a staging run between two writes and confirm row counts in every sink are unchanged after replay.
Key takeaway: foreachBatch gives you a batch DataFrame and a batchId for each micro-batch, runs between the offset log write and the commit log write, and is at-least-once by default because a crash in that window replays the same batch with the same id. Make every write idempotent with sink-level versioning such as Delta txnAppId and txnVersion, a keyed upsert, or a ledger committed in the same transaction. Cache before fan-out, let exceptions propagate, size batches at the source and watch addBatch time.