A Structured Streaming query has three parts: a source that hands out data by offset, a query plan that transforms it, and a sink that writes the result somewhere durable. Most outages that people blame on streaming are really sink problems: duplicates in a Kafka topic, an external table that shows half-written files, a dashboard that silently stops after someone deleted a checkpoint. The source and the engine can only replay data deterministically. Whether a replay becomes a duplicate is decided by the sink.

This article explains how a micro-batch reaches the sink, which sinks ship with Spark and what each guarantees according to the Spark 4.2 documentation, how the file sink's metadata log works, why the Kafka sink is at-least-once and where to deduplicate, and how to choose. It finishes with a worked pipeline that writes one stream to three sinks and a checklist you can apply to your own queries.

How a micro-batch reaches the sink

1. Plan batch Npick end offsets2. offsets/Nwrite-ahead log3. Run batch Ntasks write data4. Sink commitdriver, epoch N5. commits/Nbatch is doneCrash after 2,before 5: replay NIdempotent sinkskips or overwrites NNon-idempotent sinkwrites N twicerestartsame NorOnly the sink decides whether a replaybecomes a duplicate
The micro-batch commit sequence. A crash between the offset log and the commit log replays the same batch id with the same offsets; the sink decides whether that replay is harmless.

Everything about sink guarantees follows from one sequence that the driver runs for every micro-batch. First it decides the end offsets for batch N and writes them to offsets/N inside the checkpoint directory. That file is a write-ahead log: once it exists, batch N is defined forever. Then the batch runs; executor tasks process partitions and the sink's per-partition writers do their work. When all tasks succeed, the driver calls the sink's commit for epoch N, and only after that returns does it write commits/N.

Now consider a crash at any point between the two log writes. On restart, Spark finds offsets/N without commits/N and runs batch N again with exactly the same offsets. Stateful operators restore their state store to the version matching N-1, so the recomputed output is the same as before. The sink therefore receives the same batch id and the same rows a second time. An idempotent sink recognises the batch id and skips it, or overwrites what it wrote before. Any other sink writes it twice.

In the DataSource V2 API this contract is visible in the interfaces: each task's DataWriter returns a commit message, and the driver-side StreamingWrite.commit(epochId, messages) receives all of them with the epoch id. A custom sink achieves exactly-once only if that commit is atomic and keyed by epoch id, for example a single transaction that records the epoch alongside the data.

Sinks and output modes

The output mode decides which rows of the conceptual result table go to the sink each trigger. Append emits only rows that will never change, so aggregations need a watermark before Append is allowed. Update emits only rows that changed in this trigger. Complete re-emits the whole result table every time, which only makes sense for small aggregations. Not every sink supports every mode, and the documented guarantees differ:

SinkOutput modesDocumented guaranteeTypical use
File (parquet, json, csv, orc, text)AppendExactly-onceRaw landing zone, data lake
KafkaAppend, Update, CompleteAt-least-onceFeeding other services
foreachAppend, Update, CompleteAt-least-onceRow-at-a-time custom writes
foreachBatchAppend, Update, CompleteDepends on your codeAny batch writer, MERGE, fan-out
ConsoleAppend, Update, CompleteNoneDebugging
MemoryAppend, CompleteNone (Complete rebuilds on restart)Tests and notebooks

Table formats such as Delta Lake and Apache Iceberg plug in as additional sinks through format(...) or toTable(...); their guarantees come from their own commit protocols, covered below. Two rules from the query side matter as much as the sink column: aggregations without a watermark cannot use Append, and stream-stream joins currently support only Append.

The file sink and its metadata log

The file sink is the oldest exactly-once sink and the most misunderstood. Each task writes files with unique names into the output directory, including any partition subdirectories. Then the driver writes _spark_metadata/N in the output directory, a small log entry listing every file that batch N added. Files are only part of the dataset once they appear in that log. If the batch is replayed, the sink sees that N is already in its log and skips the batch.

That last point is the source of a nasty failure. If you delete the checkpoint but keep the output directory, the restarted query begins again at batch 0, and the file sink skips every batch whose id is not greater than the latest one in its own metadata log. The query runs, reports progress and writes nothing. Never reset a checkpoint without also choosing a new output path or deliberately clearing the old metadata log.

The exactly-once guarantee also holds only for readers that honour the log. When Spark reads the root of the output directory it finds _spark_metadata and reads only listed files. Hive, Trino, Athena or a Spark job reading a single partition subdirectory just list files, so they also see orphans left by failed task attempts and can double-count. If other engines read the output, write to a table format instead, or run a separate compaction job that publishes clean data.

Two operational costs come with the design. Every micro-batch writes at least one file per output partition, so a 10-second trigger with 200 shuffle partitions produces a flood of small files; coalesce or repartition by the partition column before the sink, and lengthen the trigger. And the metadata log lists every file ever written, so it grows without bound; Spark periodically compacts it, but the compacted file still lists everything. The documented retention option, a duration such as 7d, lets old entries drop out of the log; files beyond the retention are then invisible to Spark readers.

The Kafka sink

The Kafka sink expects a value column of string or binary type and optionally key, headers, topic and partition. If you set the topic option it overrides any topic column. A null or missing partition lets the producer's partitioner decide, so set the key deliberately: it determines ordering per key downstream.

The sink is documented as at-least-once, and duplicates come from two places. A producer may retry a send that the broker already wrote but did not acknowledge. And a replayed micro-batch resends every record of batch N. Producer-level settings passed with the kafka. prefix can reduce the first kind but cannot remove the second, because the replay is a new send of the same data. The Spark documentation's advice is the right one: include a unique key, such as an event id plus the logical version, and deduplicate when reading.

alerts = (scored
    .selectExpr("CAST(account_id AS STRING) AS key",
                "to_json(struct(event_id, account_id, score, ts)) AS value")
    .writeStream
    .format("kafka")
    .option("kafka.bootstrap.servers", "broker1:9092,broker2:9092")
    .option("topic", "fraud-alerts")
    .option("checkpointLocation", "s3://chk/fraud-alerts/")
    .outputMode("append")
    .start())

Consumers of fraud-alerts must treat event_id as the idempotency key: an upsert into a keyed store, a dedup window in a downstream stream, or a database unique constraint. Producers are cached and reused across tasks with the same configuration, so credential rotation can take as long as the producer cache timeout, ten minutes by default, to apply. For how the Kafka source side tracks offsets, see Structured Streaming with Kafka.

Table sinks: Delta, Iceberg and toTable

Delta Lake's streaming sink records the query id and batch id inside the same transaction-log commit that adds the data files. A replayed batch finds that its id is already committed and becomes a no-op, which gives exactly-once for appends without any code from you, and readers of the table see only committed snapshots. The details of the log are in Delta Lake architecture. Apache Iceberg also supports streaming appends through Spark, committing a snapshot per micro-batch; check how your Iceberg version handles replayed epochs before relying on it.

Since Spark 3.1 you can write writeStream.toTable("catalog.db.table"), which resolves the table through the catalog instead of a path. Use it: it keeps table names in code and paths in the catalog, and the catalog decides the format. Table sinks also fix the small-file problem properly, because compaction rewrites files and commits the result atomically, so readers never see a half-compacted state.

foreach, foreachBatch, console and memory

foreach gives you a ForeachWriter with open(partitionId, epochId), process(row) and close(error). It is at-least-once by construction; use the partition and epoch ids to make writes idempotent, and keep connection setup in open because it runs once per partition per epoch.

foreachBatch hands you each micro-batch as an ordinary DataFrame plus its batch id, so any batch writer works, including JDBC, MERGE into a table and writes to several destinations. Its guarantee is exactly what your function makes it. The idempotency patterns, the cost of multiple actions on one batch and caching are covered in foreachBatch in depth. It does not work with the experimental continuous trigger.

The console and memory sinks keep nothing durable. The memory sink collects results on the driver as a table named after the query, which is useful in tests and dangerous anywhere else because it grows driver memory.

Worked example: one stream, three sinks

Take a payments stream on Kafka. You want three outputs: every event landed in a Delta bronze table, per-merchant one-minute totals for a dashboard, and high-risk events pushed to an alerts topic. The rule that makes this safe is one query per sink, each with its own checkpoint, so each sink has its own commit sequence and failures are isolated.

from pyspark.sql import functions as F

events = (spark.readStream.format("kafka")
    .option("kafka.bootstrap.servers", "broker1:9092")
    .option("subscribe", "payments")
    .load()
    .select(F.from_json(F.col("value").cast("string"), SCHEMA).alias("e"))
    .select("e.*"))

# 1. Bronze landing: exactly-once through the Delta transaction log.
bronze = (events.writeStream
    .option("checkpointLocation", "s3://chk/payments-bronze/")
    .trigger(processingTime="1 minute")
    .toTable("lake.bronze.payments"))

# 2. One-minute totals: watermark first, so Append can emit closed windows.
totals = (events
    .withWatermark("ts", "10 minutes")
    .groupBy(F.window("ts", "1 minute"), "merchant_id")
    .agg(F.sum("amount").alias("total"), F.count("*").alias("n"))
    .writeStream
    .outputMode("append")
    .option("checkpointLocation", "s3://chk/payments-totals/")
    .toTable("lake.silver.merchant_minute_totals"))

# 3. Alerts: at-least-once; consumers deduplicate on event_id.
alerts = (events.where("risk_score > 0.9")
    .selectExpr("CAST(merchant_id AS STRING) AS key",
                "to_json(struct(event_id, merchant_id, amount, risk_score, ts)) AS value")
    .writeStream.format("kafka")
    .option("kafka.bootstrap.servers", "broker1:9092")
    .option("topic", "payment-alerts")
    .option("checkpointLocation", "s3://chk/payment-alerts/")
    .start())

Walk through a failure. The cluster dies after the alerts query wrote offsets/41 and sent half of batch 41's records. On restart, the bronze query replays its own pending batch, Delta finds the batch id already committed or not and acts accordingly, and the table is correct. The totals query restores state for batch 40 and recomputes 41; Append mode only emits windows the watermark has closed, and Delta deduplicates the replayed batch. The alerts query resends all of batch 41, so some alerts arrive twice with the same event_id, which the paging service ignores because it upserts on that key. Each sink behaved exactly as its guarantee says, and the system as a whole is correct because the one at-least-once sink has an idempotent consumer.

Note what the example avoids: three queries reading Kafka means three consumers of the topic, which costs broker bandwidth. If that matters, a single foreachBatch query can fan out to all three, at the price of coupling their failures and writing the idempotency logic yourself. For the totals query's state and watermark behaviour, see streaming state and watermarks.

Failure modes

  • Checkpoint deleted, output kept. File sink skips every batch up to its old latest id and writes nothing. Pair checkpoint and output directory lifetimes.
  • Two queries, one checkpoint. They overwrite each other's offset and commit logs. One checkpoint per query, always.
  • External readers see duplicates. Non-Spark engines read file-sink orphans. Use a table format or publish through a compaction job.
  • Small files. Short triggers times many partitions. Repartition before the sink, lengthen the trigger, compact tables.
  • Kafka duplicates treated as a bug. They are the documented contract. Add an event id and deduplicate downstream.
  • Complete mode on a growing aggregate. Every trigger rewrites the full result and state never shrinks. Use Update or a watermark with Append.
  • Changing the sink under an existing checkpoint. Some changes are allowed, but switching sink type or output path under an old checkpoint invites skipped or repeated data. Start a new checkpoint and backfill deliberately.

Choosing a sink

Choose by who reads the output and what a duplicate costs them. If analysts and other engines read it, use a table sink such as Delta or Iceberg; you get atomic visibility and compaction. If Spark alone reads a landing directory and you want no table format, the file sink is fine provided you never reset its checkpoint independently. If other services consume events, use Kafka and design consumers to be idempotent. If you must write to a database or call an API, use foreachBatch and make the write idempotent on batch id or business key. Reserve console and memory for development. For the wider architecture these choices sit in, see Structured Streaming architecture.

What to do next

  1. List every streaming query you run, with its sink, output mode and checkpoint path, and confirm that no two queries share a checkpoint.
  2. Mark each sink's documented guarantee next to it, and for every at-least-once sink name the downstream key that makes duplicates harmless.
  3. For each file sink, find out who reads the directory; move any directory read by non-Spark engines to a table format.
  4. Kill one query mid-batch in a staging environment and verify that the restarted output contains no gaps and no unexplained duplicates.
  5. Measure files written per hour per sink and set the trigger interval and pre-sink partitioning so file sizes are reasonable.
  6. Write a runbook rule: a checkpoint is never deleted without a decision about the sink's output and metadata.
Key takeaway: Spark replays a failed micro-batch with the same batch id and the same offsets, so the sink alone decides whether a replay becomes a duplicate. The file sink and Delta are exactly-once because they record committed batch ids, Kafka and foreach are at-least-once, and foreachBatch is whatever your code makes it. Give every query its own checkpoint, never reset a checkpoint without dealing with the file sink's metadata log, use table formats when other engines read the output, and put an idempotency key on every at-least-once stream.