skip to content

How do you assert consumer offset and commit behavior in an integration test, and why does it matter for correctness?

level: seniorimportance: should knowfreq 40%

answer

  1. committed offset = last processed + 1
  2. lag = endOffsets - committed; assert 0
  3. AdminClient.listConsumerGroupOffsets / consumer.committed()
  4. auto-commit -> loss risk; manual -> at-least-once
  5. failure test: offset must NOT advance past unprocessed

basics

~20 s

Use an admin/consumer API to read committed offsets and end offsets for the group and partitions, then assert the committed position advanced as expected. It matters because wrong commit timing causes message loss (commit too early) or reprocessing (commit too late).

solid answer

~50 s

Offset assertions verify your consumer committed the right position after processing. In an integration test you query committed offsets — e.g. AdminClient.listConsumerGroupOffsets(groupId) or a probe KafkaConsumer's committed(partitions) — and compare against endOffsets(partitions) to compute lag, asserting lag is zero (fully caught up) or at the expected position. With Spring you can read offsets via the AdminClient or assert through a ConsumerSeekAware/offset listener. This matters because commit timing is a correctness property: with enable.auto.commit=true offsets advance on a timer regardless of whether processing succeeded, risking loss on crash; with manual commit (ackMode MANUAL/MANUAL_IMMEDIATE or commitSync after processing) you guarantee at-least-once. A good test sends N records, awaits processing, then asserts committed offset == last consumed + 1 and lag == 0 — and for failure cases asserts the offset did NOT advance past an unprocessed record, proving redelivery.

go deeper

for a junior

Knows a committed offset marks where the consumer resumes and that commit timing affects loss vs duplicates.

for a middle

Can read committed vs end offsets and assert lag is zero after processing.

for a senior

Writes failure-path tests proving at-least-once (offset doesn't advance past unprocessed) and reasons about auto vs manual commit and the +1 semantics.

for a principal

Defines delivery-guarantee testing policy across services, ensuring commit semantics are asserted, not assumed, and that EOS/transaction paths have their own coverage.

**Background — offsets and commits**: Each record in a partition has an **offset** (its position). A consumer group tracks, per partition, a **committed offset** = the next offset it will read after a restart/rebalance. 'Committing' an offset of N means 'I've processed everything up to N-1; resume at N.' Getting commit timing wrong is a correctness bug: - **Commit too early** (before processing finishes) → if the consumer crashes mid-processing, those records are skipped on restart = **message loss**. - **Commit too late / not at all** → on restart the consumer re-reads already-processed records = **duplicates / reprocessing** (at-least-once). This is why tests should assert offset/commit behavior, not just 'did the message arrive.' **How to read offsets in a test**: 1. **AdminClient**: `admin.listConsumerGroupOffsets(groupId).partitionsToOffsetAndMetadata().get()` returns the committed `OffsetAndMetadata` per `TopicPartition` for the group. This is the cleanest, listener-agnostic way. 2. **A probe KafkaConsumer**: create a consumer in the **same group**, `consumer.committed(Set<TopicPartition>)` returns committed offsets; `consumer.endOffsets(partitions)` returns the log-end offset. **Lag = endOffset − committedOffset**; assert it's 0 to prove the group fully consumed and committed. 3. **Spring Kafka**: you can inject an `AdminClient`/`KafkaAdmin`, or use the listener container's offset position; `MANUAL`/`MANUAL_IMMEDIATE` ack modes let you control exactly when `acknowledge()` commits, which you then assert. **A robust offset test pattern**: ``` // 1. produce N records to the topic // 2. await (Awaitility) until the listener has processed all N // 3. assert committed offset for each partition == last consumed offset + 1 // 4. assert lag (endOffset - committed) == 0 ``` Note committed offset is **last processed offset + 1** (the next position), a classic off-by-one to get right in assertions. **Testing failure semantics** (the senior part): to prove **at-least-once**, make processing throw on a specific record, then assert the committed offset did **not** advance past that record, and that on the next poll/redelivery it's reprocessed. To catch a **loss** bug, configure auto-commit + a crash and show records were skipped — then fix by switching to manual commit after processing. This turns an abstract delivery-guarantee claim into an executable assertion. **Why integration, not unit**: real committed offsets live on the broker (in the `__consumer_offsets` topic), so only `@EmbeddedKafka`/Testcontainers exercise true commit behavior. `MockConsumer.committed(...)` can simulate offsets for unit-level logic, but it doesn't run the real commit/rebalance path — use it for logic, the broker for the actual guarantee. **Common pitfalls**: asserting the offset equals the last record's offset instead of +1; reading offsets before the async commit has flushed (wrap in Awaitility); using auto-commit in a test and wondering why offsets advance even when processing failed; and ignoring per-partition offsets by assuming a single partition.

  • How would you write a test proving at-least-once delivery on a processing failure?
    Make the listener throw on record X, await, then assert the committed offset did not advance past X (e.g. equals X, not X+1), and that on redelivery X is reprocessed. This proves no commit happened for unprocessed work, so the record is retried, not lost.
  • Why is committed offset 'last processed + 1' and not the last record's offset?
    A committed offset is the next position to read on restart. If you processed offset 9, you commit 10 so the consumer resumes at 10. Asserting against 9 is a common off-by-one error.
  • Why can enable.auto.commit=true cause message loss?
    Auto-commit advances offsets on a timer (auto.commit.interval.ms) based on what was polled, not on whether processing succeeded. If the consumer crashes after the commit but before finishing processing, those records are skipped on restart.

saying these in an interview costs you the question

  • Asserting committed offset equals the last record's offset instead of last + 1.
  • Reading committed offsets without awaiting the async commit to flush.
  • Using enable.auto.commit=true and claiming it guarantees at-least-once.
  • Claiming MockConsumer exercises real broker commit/rebalance — it only simulates offsets.

context