beginner 3 min answer Multiple choice

A consumer reads a record, writes a row to a database, then commits its offset. An engineer swaps the last two lines so the offset is committed first, reasoning that this way a record can never be processed twice. The consumer is killed roughly twice a week by deploys and node rotations. Which ordering should be the default?

kafkaoffsetsat-least-onceidempotencyconsumers
Pick one
Show the full answer Hide the answer

What is being tested

Whether you know that the two orderings have names, and that the choice between them is a choice about which failure you are willing to have. The offset store and your database are two separate systems with two separate commits. No ordering of those two commits removes the window between them. Ordering only decides which side of the window breaks.

The mechanism

Commit first, then process, is at-most-once. The consumer crashes after the offset moves and before the row is written. On restart it reads from the new offset, so the record is never processed and nothing anywhere records that it was skipped. There is no error, no retry and no dead letter. The only evidence is a number that is quietly too low.

Process first, then commit, is at-least-once. The crash happens after the write and before the commit, so the record is re-delivered and the row is written twice. That is a visible, deduplicable defect.

Put the crash rate against it. At 2000 records per second with a poll batch of 500 and a commit every 5 seconds, each unclean shutdown exposes up to 10000 records. Twice a week, commit-first silently drops up to 20000 rows a week; commit-after replays the same 20000 and a unique constraint absorbs them.

Duplicates can be defused at the sink; silent loss cannot be recovered, because the offset has already moved past the evidence. That asymmetry, not elegance, is why at-least-once is the default everywhere.

Making the write idempotent is the other half and it is cheap: a natural or derived key the producer already has (order_id plus an event sequence), a UNIQUE index on it, and an upsert or INSERT ... ON CONFLICT DO NOTHING. The idempotency key must come from the event, never from the consumer's clock or a generated id, or the second attempt looks like a different record.

Why the other options fail

  • Commit first. True that it prevents duplicates. It prevents them by converting them into undetectable loss, which is the worse of the two outcomes for anything that is counted or charged.
  • One distributed transaction. Offset stores are generally not XA participants, and even where a two-phase commit is available it buys a coordinator, a prepared-transaction state to recover and a new stall mode when the coordinator is slow. The achievable equivalent is an atomic commit of state and offsets inside one engine, or writing the offset into the same database row as the data so one commit covers both.
  • Synchronous commit per record. It shrinks the window to one record and does not close it, so the semantics are unchanged while throughput collapses: a per-record network round trip to the offset store caps a consumer at a few thousand records per second and inflates broker load.

When this is the wrong answer

At-most-once is the right default for best-effort telemetry where a lost sample is noise and a duplicate is a skewed counter - client-side metrics, sampled traces, heartbeat pings. Choose it deliberately, write down that the pipeline is lossy, and never let a billing or ledger stream inherit the setting because it was convenient for metrics.