Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. foreachBatch

PySpark · Streaming & Newer Spark

foreachBatch

Hardpyspark-88
structured-streamingforeachbatchmergeidempotency

Question

What is foreachBatch, and when do you need it?

Solution

foreachBatch hands you each micro-batch as an ordinary DataFrame, together with a batch id, and lets you run any batch code on it. You need it when the sink has no native streaming support, or when the write is more than a plain append.

Shape

def upsert_to_delta(batch_df, batch_id):
    batch_df.createOrReplaceTempView("updates")
    batch_df.sparkSession.sql("""
        MERGE INTO dim_customer t
        USING updates s ON t.customer_id = s.customer_id
        WHEN MATCHED THEN UPDATE SET *
        WHEN NOT MATCHED THEN INSERT *
    """)

(stream.writeStream
    .foreachBatch(upsert_to_delta)
    .option("checkpointLocation", "/chk/dim_customer")
    .start())

When you need it

  • Upserts with MERGE into Delta or Iceberg tables, since the streaming sinks only append.
  • Writing to a database through JDBC, or a sink that has no streaming connector.
  • Writing the same batch to several sinks (a table and a Kafka topic).
  • Running batch-only operations on each batch, such as expensive transformations that are not supported on streaming DataFrames.

Idempotency is on you

If the query fails after your function wrote the data but before Spark committed the batch to its log, Spark will rerun the same batch on restart, with the same batch_id and the same data. Your function runs twice. With a plain append to a database, that gives duplicates. The guarantee Spark gives you is at-least-once for this sink, unless you make the write idempotent.

Ways to do that:

  • Use MERGE or an upsert keyed on a business key, so a replay changes nothing.
  • Record the batch_id in a control table and skip a batch that has already been applied.
  • Delta Lake has write options (txnAppId and txnVersion) for idempotent writes, tied to the batch id.

Small tips

Cache the batch DataFrame if you write it to more than one place, and unpersist it at the end, so it is not recomputed for every write. Keep the function fast: its running time is the batch time, so a slow function means growing lag.

🎯 Put this concept into practice

Solidify this answer with real hands-on interview drills in the browser studio.

Open related drill →
PreviousNext