[CORE-16844] kafka/client: recover from offset_out_of_range in consumer group fetch - #31064
Conversation
9ca82a8 to
2293be2
Compare
940d6dd to
eb5dd14
Compare
There was a problem hiding this comment.
Pull request overview
Fixes Pandaproxy consumer-group fetch getting stuck returning offset_out_of_range after log retention/prefix-truncation advances a partition’s log start offset beyond the consumer’s initial fetch offset.
Changes:
- Update
fetch_session::apply()to seed the tracked fetch offset fromlog_start_offsetwhen a partition returnsoffset_out_of_range. - Trigger a retry from the Kafka client consumer fetch path so Pandaproxy re-fetches using the corrected offset.
- Add unit and ducktape regression tests covering prefix-trim recovery.
Reviewed changes
Copilot reviewed 4 out of 4 changed files in this pull request and generated 1 comment.
| File | Description |
|---|---|
| tests/rptest/tests/pandaproxy_test.py | Adds a ducktape regression test that trims a topic prefix and verifies a fresh consumer-group fetch still succeeds. |
| src/v/kafka/client/test/fetch_session.cc | Adds a unit test ensuring fetch_session::apply() updates tracked offsets on offset_out_of_range. |
| src/v/kafka/client/fetch_session.cc | Implements offset correction on offset_out_of_range using log_start_offset. |
| src/v/kafka/client/consumer.cc | Throws on offset_out_of_range after applying the response to force a retry before returning to Pandaproxy. |
Retry command for Build#86958please wait until all jobs are finished before running the slash command |
CI test resultstest results on build#86958test results on build#87059
test results on build#87065
test results on build#87391
|
3a87023 to
6d5f4d6
Compare
6d5f4d6 to
a037491
Compare
|
Force-pushed: squashed the per-partition recovery into the main fix commit so the PR shows the final design directly (no intermediate coarse-recovery version). Behaviour and net diff are unchanged from the previous push — only the commit layout. Diff since the previous push (6d5f4d6 → a037491): |
a037491 to
f35db72
Compare
|
Force-push: |
pgellert
left a comment
There was a problem hiding this comment.
lgtm, the main thing to consider is the timeout=0 case, the rest are minor questions
f35db72 to
15c7015
Compare
15c7015 to
e885eef
Compare
Pandaproxy's consumer group fetch started every fresh assignment at
offset 0 and never advanced once retention moved the log start offset
past it, so every fetch returned offset_out_of_range forever despite
auto.offset.reset=earliest -- the only reset policy pandaproxy accepts.
fetch_session now exposes three named operations instead of one
overloaded apply(): apply() advances offsets from a delivered response,
discard() advances only the session epoch, and reseed() sets a
partition's offset directly. consumer::fetch() collects every broker's
response, then:
- reseeds out-of-range partitions to the broker-reported
log_start_offset. earliest is exactly the log start, and the broker
already returns it in the fetch response, so no separate ListOffsets
is needed.
- strips the out-of-range partitions from the response, since the
pandaproxy serializer rejects any partition error. No offset
advances past undelivered records, so nothing is silently skipped.
- only on a dispatch failure -- a whole broker's response missing,
where the topology may be stale -- discards the round and throws, so
the retry in client::consumer_fetch() re-fetches and refreshes
metadata.
A round whose only outcome was the reseed carries no records, so instead
of returning an empty poll and deferring the data -- which already
exists at the reseeded offset -- to the client's next poll, fetch()
repeats the round within the caller's timeout budget until it has data,
nothing was reseeded, or the budget is spent. A round that delivered
records (healthy partitions) returns immediately; an out-of-range
sibling resumes on the next poll. Recovering per partition rather than
discarding the whole round avoids re-reading the healthy partitions on
every retention/trim edge.
This mirrors the in-poll recovery of franz-go and the Java consumer
(Confluent REST proxy): a timer-bounded poll() loop that resets the
position and refetches across iterations, returning data once available
or empty when the timer expires (apache/kafka 3.6):
do {
updateAssignmentMetadataIfNeeded(timer, false); // resets position
final Fetch<K, V> fetch = pollForFetches(timer); // (re)fetches
if (!fetch.isEmpty()) { ...; return records; }
} while (timer.notExpired());
return ConsumerRecords.empty();
poll loop: https://github.com/apache/kafka/blob/3.6/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java#L1174-L1207
out-of-range detect: https://github.com/apache/kafka/blob/3.6/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractFetch.java#L654-L666
reset via ListOffsets: https://github.com/apache/kafka/blob/3.6/clients/src/main/java/org/apache/kafka/clients/consumer/internals/OffsetFetcher.java#L109-L116
Adds fetch_session unit tests for the split API.
Trim a topic's log prefix past offset 0 and assert a fresh consumer group polls its way to the records at the new log start offset, producing one record per call so each lands in its own batch and the trim offset falls on a batch boundary, as real retention does. A second test trims only one of two partitions and asserts the healthy sibling's records are delivered while the out-of-range partition recovers, covering the per-partition recovery path.
e885eef to
80417ce
Compare
|
Addressed both review comments: Diff of these changes: https://github.com/redpanda-data/redpanda/compare/f35db722ac..80417ceb16 |
|
/backport v26.1.x |
|
/backport v25.3.x |
|
/backport v25.2.x |
….x-319 [v25.3.x] [CORE-16860] kafka/client: resume consumer group fetch from committed offset Follow-up to the offset_out_of_range fix (#31064, now merged into dev): a fresh HTTP Proxy consumer instance now resumes from the group's committed offset instead of always restarting at the earliest available offset. What this adds kafka/client: resume consumer group fetch from committed offset — at the start of fetch(), seed every freshly (re)assigned partition's fetch position from the group's committed offset (committed+1; our commit stores the last consumed offset). Only initializing partitions are seeded; no OffsetFetch is issued once positioned. A per-partition OffsetFetch error surfaces (throws) instead of a silent restart at earliest. tests/pandaproxy: cover resume-from-committed end-to-end — two ducktape cases: (1) instance A commits, is removed, a fresh instance B resumes from committed+1 (with a control group that has no committed offset starting at earliest); (2) a consumer commits, retention trims the log start past the committed offset, and a fresh instance resumes from committed, hits offset_out_of_range, and recovers to the new log start. This behaves very similarly to the Confluent REST API: committing offset X resumes the group at X+1.
Pandaproxy's consumer group fetch always started a fresh partition assignment at offset 0 and never advanced once retention moved the log start offset past it, regardless of auto.offset.reset=earliest requested at consumer creation. Every fetch after that point returned offset_out_of_range forever, because the client tracked its own fetch offset and had no path to correct it — and the REST client cannot control the fetch offset itself.
fetch_session is split into three named operations: apply() advances the tracked offset from a delivered response, discard() advances only the session epoch, and reseed() sets an offset directly. On offset_out_of_range, consumer::fetch() reseeds the affected partitions to the log_start_offset the broker reports alongside the error — earliest is exactly the log start, and the broker already returns it in the fetch response, so no separate ListOffsets is needed since pandaproxy only ever allows the earliest reset policy. Healthy partitions are delivered immediately (their offsets advance) and the out-of-range ones are stripped from the response (the serializer rejects any partition error). A round whose only outcome was the reseed carries no records, so rather than return an empty poll and defer the data — which already exists at the reseeded offset — to the client's next poll, fetch() repeats the round within the caller's timeout budget until it has data, nothing was reseeded, or the budget is spent, so a single fetch returns the recovered records. Only a dispatch failure, where a whole broker's response is missing and the topology may be stale, discards the round and throws so the existing gated_retry_with_mitigation in client::consumer_fetch() re-fetches and refreshes metadata.
Adds fetch_session unit tests covering the split apply/discard/reseed API, and pandaproxy ducktape regression tests: one trims a topic's log prefix and asserts a fresh consumer group polls its way to the records at the new log start offset; one trims a single partition of a two-partition topic and asserts the healthy sibling's records are delivered while the out-of-range partition recovers; and one drains a settled consumer, then trims past its position and asserts a single fetch returns the reseeded records in the same poll.
Design notes
The HTTP Proxy accepts a single reset policy. The consumer-create handler rejects anything else with a 400 before a consumer is built:
if (req_data.auto_offset_reset != "earliest") {
throw parse::error(
parse::error_code::invalid_param, "auto.offset must be earliest");
}
src/v/pandaproxy/rest/handlers.cc (create_consumer).
Because earliest is the only reachable policy, the value is validated and logged but never forwarded to the internal kafka::client — the client doesn't need to know it. This is what keeps recovery cheap: earliest resolves to the partition's log_start_offset, which the broker already returns in every fetch response, so the client reseeds straight from that field. No separate ListOffsets is issued to resolve the reset policy. (If the proxy ever allowed latest/timestamp, that assumption breaks and the value would have to be threaded through and resolved via ListOffsets.)
The Confluent REST Proxy wraps the Java KafkaConsumer, whose poll() is a timer-bounded do/while: each iteration resets any out-of-range position — via a ListOffsets to the earliest offset — and refetches, returning the recovered records within the same poll if the request timeout allows, or an empty result (recovering on a later poll) if the timer expires first. consumer::fetch() mirrors that shape:
The observable contract matches the Confluent REST Proxy: the recovered records come back in a single fetch when the timeout allows, or empty-then-next-poll otherwise; no offset_out_of_range is ever surfaced to the client, and no already-delivered records are re-read. The one mechanical difference is offset resolution — the Java consumer issues a ListOffsets(earliest) to find the log start, we reseed straight from the log_start_offset already present in the fetch response, valid precisely because of assumption (1). An out-of-range partition sitting alongside partitions that returned records is delivered on the next poll rather than the same one: the healthy records return immediately and the reseeded partition resumes from the corrected offset, just as the Java consumer's poll() returns as soon as any partition has records.
The one case that is retried is a dispatch failure — a whole broker's response missing, where the topology may be stale. There we discard the round and throw so the retry in client::consumer_fetch() runs the update_metadata refresh a stale leader needs.
Fixes: CORE-16844
Backports Required
Release Notes
Bug Fixes