Repository navigation
fix(fetch): bound each partition's read so one slow partition can't starve the rest - #852
solace-aross wants to merge 10 commits into
Conversation
b9ded0c to
c410055
Compare
sgamelin
left a comment
There was a problem hiding this comment.
Bounding each partition's read is the right fix for this, and the description makes the trade-offs easy to follow. I checked two of its Kafka claims against 3.9.1, and both hold:
DelayedRemoteFetch.onCompleteanswers an expired remote read withNONE, no records, and the offsets the broker already holds (DelayedRemoteFetch.scala#L98-L125).- The Java consumer does not add
fetch.max.wait.msto its Fetch request timeout.
I'm requesting changes for one regression, in the first inline comment. When a partition's storage budget equals its share of the read deadline, a read that uses its whole budget is always abandoned, and everything it read is lost. With a max_wait of 5s or more (franz-go's default is 5s), this covers the first partition of every Fetch with more than one partition.
The other inline comments:
- Postgres: an abandoned read puts a busy connection back into the pool. libSQL and SlateDB are only partly covered by the bound.
- Tests: the
offset_stagedeadline and theleftcount across topics have no test. - Diagnosis: a warn on each partition is the only signal of a stall.
- A non-blocking question on the halving share.
- The CHANGELOG entry.
d9b0b1b to
06cb089
Compare
06cb089 to
93c1a8a
Compare
sgamelin
left a comment
There was a problem hiding this comment.
The margin fixes the regression I requested changes for. I checked it by mutation: with Deadlines::storage capped at cap again, a_read_that_spends_its_budget_keeps_what_it_read fails, and with the margin all 12 fetch tests pass. Approving. The points below don't block the merge.
Earlier comments
- Budget reaches the cap: addressed and pinned by a test.
- Postgres connection back in the pool: partly addressed. The doc on
beforeis now correct. #827 wraps only ListOffsets in its guard (pg.rs#L2908 on #827), so pg'sfetchstays uncovered after both PRs merge, and no issue tracks it. Could you open one before this merges? The description's "Postgres" bullet also still says "the abandoned request held its connection just as long", which contradicts the header above it. - Tests for
offset_stageandleft: addressed by the three new tests. - Diagnosis: addressed by the counter and the debug level for reads that never started. One follow-up is inline.
- Halving: your reasoning about wide requests is convincing. Leaving it.
- CHANGELOG: addressed. One wording point is inline.
New
- libSQL permit wait:
SemaphoreProxy::fetchpasses the budget on unchanged after it waits for its permit (proxy.rs#L180-L197). libSQL's defaultSemaphoremode has one permit per broker. A read that gets the permit just before its cap then runs synchronously for its whole budget, up to three quarters of its share past the cap. This isn't a regression, since main has no cap at all. Passingmax_wait.saturating_sub(start.elapsed())there would charge the wait to the budget, and it's a one-line change. - Description: a few parts predate the margin and are now stale:
- the "Storage budget" formula;
- "keeps today's timing, bounded by
max_wait", which no longer holds for the first partition from amax_waitof 3s; - the mutation table;
- the header's reference to 06cb089, which is not on the branch.
- Merge order with #826: #826 reads
offset_stagebefore the fetch, outside any deadline, and itspartition_datatakes anOffsetStagerather than anOption. Whichever PR merges second has to bring that first read underbefore, and choose what a partition answers when it stalls there.
93c1a8a to
0125c1c
Compare
|
Rebased on main after #826 and ready for another look. The three inline threads are answered and resolved. The rest of your approval notes:
Whichever of #827 and #852 merges second folds its deadline helpers onto the other's. That is about 70 lines out of |
sgamelin
left a comment
There was a problem hiding this comment.
The rebase onto #826 works well. The first offset stage read now runs under the partition's share, the budget starts once it's read, and a partition whose records stall answers with the offsets it read first rather than -1. I checked the new tests by mutation: reading the first offset stage outside before fails a_stalled_offset_stage_before_the_read_answers_unknown_offsets, and all 22 fetch and proxy tests pass at 0125c1c. The -1 answer for a partition whose offsets weren't read in time is safe for Java: at 3.9.1 the consumer updates the high watermark, log start and last stable offset only when each is non-negative (FetchCollector.java#L292-L311). Approving again. The points below don't block the merge.
Earlier comments
- Margin doc and tests: addressed. The two short-share tests pin the 50ms least margin and the zero budget.
offsetin the warn, and the labels: addressed, using #827'soperationandstage.- CHANGELOG wording: addressed.
- Permit wait in
SemaphoreProxy::fetch: addressed, with a test. - Stale description: addressed.
- Merge order with #826: settled by this rebase.
- An issue for pg
fetchreturning a busy connection to the pool: partly addressed. The description says these gaps are tracked, but I couldn't find GitHub issues for pgfetch, the synchronous libSQL statement, or SlateDB ignoring its budget. Could you open them here and link them from "Not covered", so a reader can follow them from the PR?
New
Two inline comments: a missing test for the ReadCommitted arm of covering, and libSQL's mpsc mode.
0125c1c to
3d330c0
Compare
…tarve the rest Fetch read its partitions one after another with no bound on any single storage read. A partition whose read never finished held the whole request until the client gave up (30s for Java), so every partition after it in the same request got nothing, on every retry. Each Fetch now has a read deadline of max_wait plus 5s, under the read timeouts of Java, librdkafka and franz-go. Each partition may spend half of the time left before it, the last partition all of it. A partition whose read misses its share is answered with no error and no records, so the client fetches it again, and the partitions after it are still read. dynostore, lite and postgres stop assembling batches once their storage budget is spent, so the budget left of the client's max_wait, already used up by the slow partition, would read nothing. A partition now gets whatever is left of max_wait, but never less than half of max_wait of its own, and never past its share of the read deadline. max_wait still bounds a healthy request as before. Batches read before a stall are kept, and their bytes stay counted against max_bytes. The partition count is taken again on every round of the long poll. A read is not started once its deadline has passed. The long poll measures time with Instant rather than SystemTime. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D2qdgPhMyLGZN5dLR8CVHs Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
Per .claude/rules/comments.md: a comment is a complete, capitalized sentence. Three new comments in this PR started lowercase. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
A read that spent its whole budget returned just after the partition's share of the read deadline, so the broker abandoned it and lost every batch it had assembled. The budget now ends a quarter of the share, and at least 50ms, before it. Count each abandoned read in nisshi_storage_read_deadline_exceeded by stage, log a read that never started at debug, document how far the deadline reaches into each engine, and test the offset stage deadline and the partition count across topics. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_011P97dHLhTRdJ35fMJFpYPg Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
… log the offset The `stage` attribute of nisshi_storage_read_deadline_exceeded held either where the read was (`read`, `offset_stage`) or why it was missed (`not_started`), so a records read and an offset stage read that never started counted the same. `operation` now names the Storage method that was read (`fetch` or `offset_stage`), and `stage` holds `reading` or `not_started`, the values the ListOffsets deadline uses for the same counter. The log lines carry the operation and the offset, so a stall at one offset or object can be told from a slow engine. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
The doc on Deadlines::storage said a partition that starts late still gets half of max_wait, which the margin cut short once three quarters of the share is under that half. It now says what each partition gets. The margin tests took an overshoot of 8ms with shares of 2.75s or more, so any margin passed them. A read that overshoots by 1s at a 5s max_wait now pins the quarter rule, a read after five stalls with a 172ms share pins the 50ms least margin, and a share under the least margin pins the zero budget. A stall of the offset stage before the records, and a stall of the records past the stage, each get a test, since the offset stage is now read before the records. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
A stage read that was slow but answered would otherwise leave its partition what the stage left of max_wait, rather than the half the budget promises. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
…:fetch libSQL's default Semaphore mode has one permit per broker. A read that got the permit late ran for its whole budget from then, up to three quarters of its share past the point the Fetch deadline expects it back, and was abandoned with what it had read. The budget passed on is now what is left of it after the wait. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
The entry said a partition that runs out of time answers -1 offsets, and that the broker logs every such partition at warn. The offsets are now read before the records, so a partition whose records stall answers the offsets it has, and a partition never started logs at debug. The entry names the counter's operation and stage values. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com> # Conflicts: # CHANGELOG.md
… unread partition counts Only dynostore reads nothing with a zero budget; postgres and libSQL read one record, and SlateDB ignores the budget. An offset stage read again after the records takes a fresh half of the time left, so a partition can spend up to three quarters of it. A partition left unread after a stall counts under offset_stage, since that is the first read. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
A partition that stalls after returning records, and one whose offsets re-read stalls, now run under both isolation levels and assert that the last stable offset is raised to cover the records under ReadCommitted. Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com> Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
3d330c0 to
444365b
Compare
What does this change do, and why?
Fetch reads its partitions in order, and nothing bounded a single storage read. If one partition's read never finished, the request waited until the client gave up (30s for Java), and every partition after it in the same request got nothing, on every retry. This bounds each partition's Fetch reads by a deadline of its own, in
nisshi-storage/src/service/fetch.rs.max_wait + 5s, under the Fetch read timeouts of Java (30s), librdkafka (60s +fetch.wait.max.ms) and franz-go (10s +max_wait).max(end of max_wait, now + max_wait / 2), capped atcap - margin, wherecapis the partition's share andmarginis a quarter of the share, at least 50ms, so a read that spends its whole budget returns what it assembled before the broker abandons it.max_waitbudget whilemax_waitis under 3s. From 3s the first partition gets three quarters of its share: 3.75s at franz-go's 5s default.max_waitof its own, as far as its share allows, and no budget once the share is under 50ms.max_wait. That trades latency in that case for no partition going unread, still under every client's own timeout.SemaphoreProxy::fetch(libSQL's default mode, one permit per broker) charges the wait for its permit to the budget it passes on.ErrorCode::None, the batches it had read, and the offsets read before them, with the high watermark raised to cover the batches (and the last stable offset too, underReadCommitted), so the client fetches it again. The bytes those batches took stay counted againstmax_bytes.Nonewith no records and -1 for the high watermark, last stable offset and log start offset. Java ignores negative offsets; librdkafka and franz-go report -1 watermarks for that partition until the next full answer. Kafka answers an expired remote read (DelayedRemoteFetch) withNoneand the offsets it holds; the -1 case is nisshi's own.Instantinstead ofSystemTime.timeout_atpolls its future once before checking the clock.nisshi_storage_read_deadline_exceededwithoperation(fetch,offset_stage) andstage(reading,not_started), the names and values fix: bound each ListOffsets partition read by a deadline #827 uses. A read abandoned in storage logsfetch read deadline exceededat warn with topic, partition and offset; a read never started logsfetch read not startedat debug.What clients see: the stalled partition is retried on every round, so the partitions behind it are served once per round, not once per client timeout. librdkafka moves the stalled partition to the end of its next request; Java does not, so for Java it stays first and costs up to half the read deadline each round.
Not covered by the deadline:
statement_timeout. The budget margin means pg's own budget check usually ends the read first. fix: bound each ListOffsets partition read by a deadline #827 adds a guard that takes a dropped read's connection out of the pool and cancels its statement; pg'sfetchcan use it once fix: bound each ListOffsets partition read by a deadline #827 lands.mpscmode, where a fetch already in the request channel still runs with its whole budget after Fetch abandons it. Passing anInstantdeadline throughStorage::fetchwould close this and the permit-wait accounting together.storage.metadata, which is still unbounded. Partitions are still read one after another.beforeandMissedhere matchdeadline::withinanddeadline::Missedfrom #827, and the counter has the same name, attributes and values. Whichever PR merges second drops its copy and calls the other's.How was this tested?
The unit tests in
fetch.rsuse a storage double whose reads are scripted per partition: a read answers at once, never answers, or, like dynostore and lite, reads nothing with a zero budget and returns once its budget is spent plus an overshoot. They cover a stalled partition not starving the rest, stalled partitions ending at the read deadline, a partition that stalls after one batch keeping it (under both isolation levels, asserting the raised last stable offset underReadCommitted), long-poll rounds sharing the deadline again, both budget margin rules, a zero budget under the least margin, the budget starting after the offsets, offsets stalling before and after the records (both isolation levels), two topics sharing one deadline, an unknown topic counting as read, no polling after the deadline, andSemaphoreProxy::fetchcharging its permit wait.Each rule was mutated and the tests re-run:
RequestTimedOutinstead ofNonefor stalled offsetsmax_waitReadCommittedLocally:
just fmt,just clippy,cargo nextest run -p nisshi-storage --all-features(53 passed), and the nisshi-brokerfetch,fetch_wire,produce_fetch,storage_fetch,storage_list_offsetsandlist_offsetstests on in-memory, libSQL, Postgres and SlateDB (81 passed).Checklist
git commit -s) — see CONTRIBUTING.mdjust fmt,just clippy, andjust testpass locally (targeted tests as above; CI runs the full suite)🤖 Generated with Claude Code