skip to content

How does each engine achieve exactly-once processing in a Kafka pipeline, and what is the scope of that guarantee?

level: seniorimportance: must knowfreq 70%

answer

  1. exactly_once_v2, single transaction, offsets in txn
  2. read_committed downstream
  3. Flink: barrier snapshot + 2PC sink
  4. Spark: checkpoint offsets + idempotent/txn sink
  5. EOS only within participating systems

basics

~10 s

Kafka Streams uses Kafka transactions (set processing.guarantee=exactly_once_v2) so reads, state updates, and writes commit atomically. Flink uses checkpoint barriers plus a two-phase-commit Kafka sink. Spark uses checkpointed offsets with idempotent/transactional sinks.

solid answer

~40 s

All three give exactly-once *within the Kafka boundary* — consuming, updating state, and producing become atomic — not exactly-once side effects everywhere. Kafka Streams sets processing.guarantee=exactly_once_v2, wrapping the consume-process-produce loop and state-store changelog writes in a single Kafka transaction committed via the producer's transactional API; offsets are committed to the consumer group as part of that transaction. Flink achieves it with distributed snapshots (Chandy-Lamport-style checkpoint barriers) for state, plus a TwoPhaseCommitSink (KafkaSink with DeliveryGuarantee.EXACTLY_ONCE) that pre-commits on checkpoint and finalizes on checkpoint-complete. Spark Structured Streaming checkpoints source offsets and state to durable storage and relies on the sink being idempotent or transactional (the Kafka sink supports it). The crucial caveat: exactly-once is end-to-end only if every hop participates; an external HTTP call or non-transactional DB write inside your processor is at-least-once.

go deeper

for a junior

Know that all three can do exactly-once into Kafka, and that Streams uses processing.guarantee=exactly_once_v2.

for a middle

Explain the Kafka transaction wrapping consume/state/produce, and that downstream needs read_committed.

for a senior

Compare mechanisms: Streams transactions, Flink barrier-snapshot + 2PC sink, Spark checkpoint + idempotent sink; articulate the end-to-end scope caveat.

for a principal

Design pipelines for end-to-end correctness: idempotent side effects, transaction.timeout tuning, and the latency cost of read_committed and small commit intervals.

## What 'exactly-once' actually means here Exactly-once-semantics (EOS) does **not** mean a record is physically processed once — failures cause reprocessing. It means the **observable effects** are as if each input record affected the output and state **exactly once**: no duplicates, no lost updates. In Kafka pipelines this is scoped to the **consume → process → produce** loop plus internal state. ## Kafka Streams Set `processing.guarantee=exactly_once_v2` (the `_v2` variant, post Kafka 2.5/3.0, uses a single producer per instance instead of producer-per-task and is far more scalable than the original `exactly_once`). Mechanism: 1. The app reads records, updates **state stores** (RocksDB), and produces output. 2. State-store mutations are also written to **changelog topics**. 3. All output writes, changelog writes, **and the consumer offset commit** are bundled into one **Kafka transaction** using the idempotent/transactional producer (`transactional.id`, `enable.idempotence=true`). 4. On `commit.interval.ms`, the transaction commits atomically; on failure it aborts and the records are reprocessed from the last committed offset. Consumers downstream must set `isolation.level=read_committed` to skip aborted records. Trade-off: lower `commit.interval.ms` (Streams defaults it to 100ms under EOS) trims end-to-end latency but adds transaction overhead. ## Apache Flink Flink decouples **state consistency** from **sink delivery**: - **State**: Flink periodically injects **checkpoint barriers** into the stream (asynchronous barrier snapshotting, a variant of the Chandy-Lamport algorithm). When a barrier flows through every operator, each snapshots its state to the configured **state backend** (e.g., RocksDB) with the snapshot persisted to durable storage (S3/HDFS). This gives exactly-once *state*. - **Sink**: For exactly-once *output* to Kafka, the `KafkaSink` with `DeliveryGuarantee.EXACTLY_ONCE` implements **two-phase commit**: on each checkpoint it *pre-commits* a Kafka transaction; when the checkpoint is globally complete it *commits* it. If the job dies before commit, the transaction is aborted/recovered. You must set a `transactionalIdPrefix` and tune `transaction.timeout.ms` to exceed the checkpoint interval, or transactions expire and you get data loss/duplicates. ## Spark Structured Streaming Spark uses a **checkpoint location** (a directory on durable storage) holding a **write-ahead log of offsets** and a **state store** (versioned, with HDFS/RocksDB backends). On recovery it replays from the last committed batch. Exactly-once *output* requires the sink to be **idempotent** (deterministic batch IDs) or **transactional**; the built-in Kafka sink and file sink support it, but arbitrary `foreach` sinks do not unless you make them idempotent. ## The universal caveat — end-to-end scope EOS is **transitive only across participating systems**. If your processor makes a non-transactional side effect — an HTTP POST, a write to a DB outside the transaction, sending an email — that effect is **at-least-once** and will be duplicated on reprocessing. Designing for end-to-end EOS means making side effects idempotent (dedup keys, upserts) or keeping them inside the transactional boundary. ## Edge cases - **Zombie fencing**: transactional producers fence out stale instances via epoch bumps so a paused-then-resumed old instance can't double-write. - **Transaction timeout vs checkpoint interval** (Flink) is the classic foot-gun: `transaction.timeout.ms` must be larger than the maximum expected time between checkpoints, bounded by the broker's `transaction.max.timeout.ms`. - **read_committed lag**: consumers may see higher end-to-end latency because they only read up to the last stable offset (LSO).

  • Why is exactly_once_v2 preferred over the original exactly_once in Kafka Streams?
    The original used one transactional producer per task, so producer count exploded with partitions/tasks, harming throughput and broker load. exactly_once_v2 uses a single producer per instance (thread) and commits all tasks' work together, scaling far better. It requires brokers 2.5+.
  • A Flink job to Kafka shows occasional duplicates after recovery. What's the likely misconfiguration?
    transaction.timeout.ms is too short relative to the checkpoint interval (or restart gap), so the broker aborts/expires the pre-committed transaction before Flink commits it. Increase it (and broker transaction.max.timeout.ms) to comfortably exceed checkpoint interval plus recovery time.
  • Does exactly-once protect a database write inside your processor?
    No, unless that write is part of the same transaction or made idempotent. Only the Kafka consume/state/produce loop is atomic; external side effects are at-least-once and need dedup/upsert semantics.

saying these in an interview costs you the question

  • Claiming exactly-once means a record is literally processed only once — it means effects are as-if-once; reprocessing still happens on failure.
  • Saying exactly-once covers arbitrary external side effects (HTTP/DB) automatically.
  • Forgetting downstream consumers must use isolation.level=read_committed to actually get EOS.
  • Confusing Flink's checkpoint barriers (state) with its sink's two-phase commit (output) — they are separate mechanisms.

context