skip to content

How does manual acknowledgement work in a Reactor Kafka stream, and what is the role of ReceiverOffset.acknowledge() vs commit()?

level: seniorimportance: must knowfreq 45%

answer

  1. ReceiverOffset per ReceiverRecord
  2. acknowledge() = mark eligible, batched commit
  3. commit() = immediate Mono<Void> checkpoint
  4. commitInterval + commitBatchSize
  5. ack after success, concatMap for order

basics

~20 s

Each ReceiverRecord exposes a ReceiverOffset. Call acknowledge() after you finish processing a record to mark its offset eligible for committing; Reactor Kafka periodically commits acknowledged offsets in batches. commit() forces an immediate, explicit commit and returns a Mono.

solid answer

~50 s

In Reactor Kafka you typically disable enable.auto.commit and use manual offset handling so a record is committed only after successful processing (at-least-once). Every ReceiverRecord carries a ReceiverOffset handle. After processing, you call receiverOffset.acknowledge(); this doesn't commit immediately — it marks the offset as processed, and Reactor Kafka commits acknowledged offsets in batches according to commitInterval (time) and commitBatchSize (count) in ReceiverOptions, committing the highest contiguous acknowledged offset per partition. If you need a synchronous, deterministic commit at a specific point, call receiverOffset.commit(), which returns a Mono<Void> you compose into the stream and that completes when the broker confirms. The key rule: acknowledge/commit only after side effects succeed, and preserve per-partition order (concatMap, not flatMap) so you never acknowledge offset N+1 before N's work is durable, which would risk silent message loss on rebalance or crash.

code

java · 9 lines
java
KafkaReceiver.create(receiverOptions)   // enable.auto.commit=false
    .receive()
    .concatMap(record ->
        process(record.value())              // returns Mono<Void>
            .doOnSuccess(v -> record.receiverOffset().acknowledge())
            .onErrorResume(e -> { log.error("fail, will redeliver", e); return Mono.empty(); })
    )
    .subscribe();
// commitInterval/commitBatchSize in receiverOptions control batched commit of acked offsets

go deeper

for a junior

Know that you ack a record after handling it so its offset gets committed.

for a middle

Distinguish acknowledge() (batched, eligibility) from commit() (immediate Mono) and why ack-after-success.

for a senior

Reason about commitInterval/commitBatchSize, contiguous-offset semantics, ordering with concatMap, rebalance reprocessing.

for a principal

Set org-wide delivery-guarantee patterns: idempotent consumers, when to use transactions vs at-least-once, checkpoint placement.

## The delivery-guarantee problem A Kafka consumer tracks progress via **committed offsets** per (topic, partition, group). On restart or rebalance, consumption resumes from the last committed offset. *When* you commit decides the guarantee: - Commit **before** processing → at-most-once (crash = lost message). - Commit **after** processing → at-least-once (crash = possible reprocess/duplicate). **Auto-commit** (`enable.auto.commit=true`) commits on a timer regardless of whether processing finished — unsafe for at-least-once. So reactive consumers usually set it false and commit manually. ## ReceiverOffset Each `ReceiverRecord<K,V>` you get from `receiver.receive()` exposes `receiverOffset()` returning a **ReceiverOffset**, which knows its topic-partition and offset and offers two methods: ### acknowledge() - Marks this offset as **processed/eligible to commit**. It does NOT commit synchronously. - Reactor Kafka's commit scheduler then commits acknowledged offsets in **batches**, governed by `ReceiverOptions`: - `commitInterval` (e.g. 5s) — commit at least this often. - `commitBatchSize` (e.g. 100) — commit once this many acks accumulate. - It commits the **highest contiguous acknowledged offset** per partition, so out-of-order acknowledgement can stall the committed position until the gap is filled. - Efficient: amortizes commit RPCs across many records. ### commit() - Returns a `Mono<Void>` that performs an **immediate, explicit** offset commit and completes when the broker confirms. - Use when you need a hard checkpoint (e.g. after a batch boundary, before a critical external effect, or at shutdown). You compose it into the stream: `.concatMap(rec -> process(rec).then(rec.receiverOffset().commit()))`. ## Ordering matters Because committing offset N implies N-1 is done, you must not acknowledge later offsets before earlier ones are durably processed. Use **concatMap** (sequential, order-preserving) for processing within a partition. `flatMap` runs concurrently and interleaves completions, which can acknowledge a high offset while a lower one's side effect hasn't happened — a classic source of silent loss after a rebalance. ## Out-of-order / partial failure - If processing fails, do NOT acknowledge — let the stream error or retry (`retryWhen`) so the offset is re-delivered. - On rebalance, only committed offsets survive; in-flight unacknowledged work is reprocessed by the new owner (hence at-least-once + idempotent consumers). ## Auto-ack convenience - `receiver.receiveAutoAck()` returns `Flux<Flux<ConsumerRecord>>` and auto-acknowledges each inner batch after it's consumed — simpler but gives up fine-grained control. - `receiveAtmostOnce()` commits before delivery for at-most-once. ## Summary rule Acknowledge/commit AFTER success, preserve order, prefer batched `acknowledge()` for throughput and explicit `commit()` (the Mono) for deterministic checkpoints.

  • What governs how often acknowledge()'d offsets are actually committed?
    commitInterval (time-based) and commitBatchSize (count-based) in ReceiverOptions; a commit fires when either threshold is hit, committing the highest contiguous acknowledged offset per partition.
  • Why can using flatMap for per-record processing cause message loss?
    flatMap processes records concurrently and completes out of order, so a higher offset can be acknowledged before a lower offset's side effect is durable; a crash/rebalance then skips the un-done lower offset. concatMap preserves order.
  • Why disable enable.auto.commit in reactive consumers?
    Auto-commit fires on a timer independent of processing success, breaking at-least-once. Manual acknowledge()/commit() ties the committed offset to completed processing.

saying these in an interview costs you the question

  • Saying acknowledge() commits to the broker immediately (it only marks eligibility; commit is batched).
  • Acknowledging before processing side effects complete (turns at-least-once into at-most-once / loss).
  • Using flatMap for ordered per-partition processing and then acknowledging out of order.
  • Leaving enable.auto.commit=true while claiming at-least-once semantics.

context