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
MERGEinto 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
MERGEor an upsert keyed on a business key, so a replay changes nothing. - Record the
batch_idin a control table and skip a batch that has already been applied. - Delta Lake has write options (
txnAppIdandtxnVersion) 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.