concept

Offset Commit Ordering

also called Commit-Then-Process, Process-Then-Commit

Where a consumer's offset commit sits relative to the side effect it is recording - the single line of code that decides whether a crash loses records silently or duplicates them visibly.

kafkaoffsetsat-least-onceidempotencyconsumers

A consumer does three things per record: read it, act on it, and record that it is done. The last two are writes to different systems, and a process can die between them. Which of the two runs first is not a style question, it is the choice of failure mode, and it is usually made by whoever typed the loop.

Put a crash rate against it. A consumer handling 2000 records a second with a commit every 5 seconds has up to 10000 records in flight behind each commit. Deploys, node rotations and autoscaler evictions give most teams one or two unclean shutdowns a week. Commit-then-process loses up to 10000 records per event with no error raised. Process-then-commit replays the same 10000 and a unique constraint absorbs them.

The library's framing is usually "at-most-once versus at-least-once". Those are names for the two orderings, and the reason the industry standardised on the second is that a duplicate is evidence and a silent skip is not.

Why it matters

The commit is a claim about work that happened somewhere else. Because there is no transaction spanning the offset store and the sink, no ordering removes the window - it only decides which side of the window is wrong.

That asymmetry is what settles the default. A duplicate row can be removed, ignored or absorbed by an idempotency key at any time after the fact. A record that was skipped is gone: the offset has advanced past it, no retry exists, nothing landed in a dead-letter topic, and the only trace is a total that is slightly low. Teams discover at-most-once pipelines months later during a reconciliation, and cannot say which records are missing.

The second reason it matters is that the ordering is often chosen accidentally. Kafka consumers with enable.auto.commit enabled commit the offsets of the previous batch on a timer (5 seconds by default) from inside poll(), with no knowledge of whether your handler finished. A handler that is slow, or that hands records to a thread pool, is running under at-most-once semantics without anyone having decided so.

Implementation patterns

  • Process, then commit, with an idempotent sink. The default. Derive the key from the event (order_id plus an event sequence), put a UNIQUE index on it, and write with an upsert or INSERT ... ON CONFLICT DO NOTHING.
  • Never derive the idempotency key from the consumer. A generated UUID or now() makes the retry look like a new record, which is the commonest way a correctly ordered consumer still produces duplicates.
  • Store the offset in the sink. If the sink is transactional, write the data and the offset in one transaction and ignore the broker's offset store on restart. One commit, no window.
  • Atomic state-and-offset commits inside the engine. A stream processor's exactly-once mode snapshots operator state and source offsets together, which closes the window for state. It says nothing about an external sink that cannot join that commit.
  • Turn auto-commit off in anything that counts, and commit explicitly after the batch's writes have been acknowledged.

Industry example

Kafka's own client documentation describes auto-commit as committing on the next poll(), which is why the long-standing production advice is to disable it for any consumer whose work must be durable. The pattern that replaced it in practice is the pairing of at-least-once delivery with a deduplicating sink: transactional producers and idempotent writes, with a reconciliation job as the backstop. The guarantee is end-to-end effect-once, assembled from an ordering choice and a constraint, not switched on by a flag.

Failure scenarios

  • The silent shortfall. Commit-first plus weekly restarts gives a counter that is a fraction of a percent low forever. No alert can see it, because nothing failed.
  • The accidental at-most-once. Auto-commit left on, records handed to an executor, a crash during a deploy - and offsets have moved past work that never ran.
  • The ordered duplicate. Process-then-commit with a non-idempotent insert produces duplicate rows at exactly the rate of unclean shutdowns, which is low enough to look like a data-quality mystery rather than a design defect.
  • The partial batch. Commit after a batch where record 300 of 500 failed, with the failure swallowed by a broad catch. The offset advances over the gap.

Trade-offs

Choose Gains Pays
Commit after + idempotent sink losses are impossible; defects are visible a key design and an index on every sink
Commit before no duplicate work after a crash undetectable loss proportional to restart rate
Offset inside the sink transaction one commit; no window at all only works where the sink is transactional

When not to use it

At-most-once is the right default for best-effort telemetry: sampled traces, client-side metrics, heartbeats, where a dropped sample is noise and a duplicated one skews an aggregate. Choose it explicitly, write down that the stream is lossy, and keep that configuration away from anything a ledger reads.

The other case for stopping short is a sink that genuinely cannot be made idempotent - sending an email, calling a third-party payment API. Then the ordering choice is not enough and you need an outbox with a claim check, which is a different pattern and more work.

Interview question

Q: A junior engineer moves the offset commit above the database write and says it prevents duplicates. Explain what they have actually changed, and what you would require before approving either ordering.

What a strong answer covers: the two orderings are at-most-once and at-least-once; the window between the two systems cannot be removed by ordering; loss is undetectable and duplicates are not, so at-least-once is the default; the requirement is an event-derived idempotency key with a uniqueness constraint, not a code review promise; auto-commit must be off; and the one case where at-most-once is correct is cheap telemetry, declared in writing.

Quick check

Quiz: A consumer commits offsets before processing. A deploy kills it mid-batch. What is the symptom? Nothing - no error, no retry, no dead letter. Only a total that is quietly low, found weeks later in a reconciliation.

Flashcard: Which ordering is the safe default and why? - Process then commit, with an event-derived idempotency key on the sink: duplicates can be absorbed, silent loss cannot be recovered.