beginner 2 min answer Multiple choice

An hourly batch job writes each hour's aggregate into a dt/hour partition and is safe to rerun after any crash. The team reimplements the same logic as a streaming job writing to the same table. The streaming job is at-least-once and now produces duplicate rows after every restart. What made the batch version's retries safe?

streamingbatchidempotencysinksat-least-once
Pick one
Show the full answer Hide the answer

The mechanism

The batch job's output has an address derived from its input: everything the 14:00 run produces belongs in dt=2026-10-03/hour=14, and nothing else writes there. A rerun deletes that directory and writes it again. The write is idempotent because the output is addressable, not because the engine is clever. Crash halfway through, rerun, and the only trace is a slightly later file timestamp.

A streaming job has no such address. It processes record 4,812,993 and emits a row. There is no name for "the set of rows this run produced", so when the job restarts from its last committed offset and re-emits the records after it, the sink has no way to recognise them as the same rows. At-least-once delivery plus an append-only sink equals duplicate rows, and the duplicate count equals the records between the last commit and the crash.

Why the other options fail

  • "No state between runs." True and irrelevant. A stateless streaming job that filters and forwards has the same duplicate problem, because the issue is on the output side. Statelessness changes recovery time, not write safety.
  • "Bounded input so it can wait for completeness." A real difference, and it buys something else: the batch number never changes after publication, while a streaming aggregate is provisional. It does nothing for retry safety. A batch job over a complete hour would still double-count if it appended instead of overwriting.
  • "Batch commits transactionally and streams cannot." Backwards. Flink and Kafka Streams both commit state and offsets atomically, and transactional sinks exist. Many batch writers are plain file writes with no transaction anywhere. The property doing the work is the overwrite semantics of a partition, which a stream can also have if you give it one.

How to make the stream safe

  1. Give every output row a deterministic key derived from the input: source record id, or (entity, window start). Then the sink can upsert instead of append, and a re-emitted row overwrites itself.
  2. If the sink has no key, write to a staging location named by the window and swap by pointer once the window closes. This is the batch trick applied to a stream.
  3. If neither is possible, pair offsets with output inside one transaction, which is what exactly-once sink connectors do, and accept that output is now only visible at commit boundaries.

When not to move this to a stream at all

Freshness is the only reason to accept this work. If the consumers of that table look at it a few times a day and tolerate data up to an hour old, the hourly job already meets the requirement and hands you idempotency, reproducibility and a trivial backfill (rerun the date range) at no cost. The streaming version costs an always-on job, a sink key design, a checkpoint store and a late-data policy. Spend that only when somebody can state a decision that gets worse at 60 minutes of staleness and better at 60 seconds.