This is an automated email from the ASF dual-hosted git repository. hubcio pushed a commit to branch benchmark/deferred-polling-task-per-partition in repository https://gitbox.apache.org/repos/asf/iggy.git
commit 0da8329e42bcfb14875bd56e769eac2bd900272c Author: Hubert Gruszecki <[email protected]> AuthorDate: Thu Sep 17 19:11:46 2026 +0200 fix(partitions): serve the prefix a byte-limited disk read already holds A disk walk that spent its byte budget discarded every batch it had matched and returned InvalidSizeBytes. The consumer polled the same offset again, hit the same budget, and was refused again, so a partition whose reads came off disk stopped making progress. Deferred polls take the whole failure: the immediate path never surfaces this error, and the Java client does not retry it, so the driver drops the assignment. Return the matched prefix instead, and refuse only when the budget was spent before the first match, which means a single batch cannot be represented. A short reply is a valid answer: the consumer advances and asks for the rest. The walk now reports a spent budget separately from an I/O fault, so a prefix is not mistaken for unreadable data, and the result carries the flag so the caller does not splice the resident tail onto it. That splice would skip the offsets still on disk and spend bytes the read never reserved. --- core/partitions/src/poll_plan.rs | 113 ++++++++++++++++++++++++++++++++++----- 1 file changed, 99 insertions(+), 14 deletions(-) diff --git a/core/partitions/src/poll_plan.rs b/core/partitions/src/poll_plan.rs index 525e41df6..a6fbd9cec 100644 --- a/core/partitions/src/poll_plan.rs +++ b/core/partitions/src/poll_plan.rs @@ -428,9 +428,17 @@ impl PollPlan { last_matching_offset, matched, faulted, + limited, } => { read_error = faulted.then_some(IggyError::CannotReadMessage); - let remaining = query.count().saturating_sub(matched); + // A budget-limited prefix stops here. Splicing the journal + // onto it would both skip the offsets still on disk and + // spend more bytes than the read reserved. + let remaining = if limited { + 0 + } else { + query.count().saturating_sub(matched) + }; let continuation = last_matching_offset .and_then(|last_offset| { resident_tail.straddle_continuation( @@ -520,12 +528,16 @@ pub enum PollTier { /// a `Faulted` result, or it would splice the next resident op over the /// unreadable run and silently skip live messages. pub enum DiskReadOutcome { - /// Walk produced matches (possibly a partial prefix if a fault stopped it). + /// Walk produced matches (possibly a partial prefix if a fault or the byte + /// budget stopped it). Matched { fragments: PollFragments<4096>, last_matching_offset: Option<u64>, matched: u32, faulted: bool, + /// The budget stopped the walk. Later offsets exist but are not in this + /// result, so nothing may be appended to it. + limited: bool, }, /// Walk completed with no fault and matched nothing. The query offset is /// below disk retention too, so the caller may serve the journal forward @@ -557,6 +569,9 @@ enum SegmentWalk { /// Fail-closed: the segment may hold present-but-unreadable or corrupt /// data, so no later segment may be served over it. Faulted, + /// The byte budget is spent. Whatever matched is a contiguous prefix and is + /// serveable, but the walk may not continue. + Limited, } /// The state one disk walk carries across its segments. @@ -735,15 +750,17 @@ impl DiskReadPlan { break; }; - if matches!( - self.walk_segment(&file, query, count, persisted, &mut walk) - .await, - SegmentWalk::Faulted - ) { - faulted = true; - break; + match self + .walk_segment(&file, query, count, persisted, &mut walk) + .await + { + SegmentWalk::Done => walk.position = 0, + SegmentWalk::Limited => break, + SegmentWalk::Faulted => { + faulted = true; + break; + } } - walk.position = 0; } // The three ratios a read-sizing change is judged on: bytes asked of @@ -762,7 +779,10 @@ impl DiskReadPlan { "disk poll read accounting" ); - if walk.limited { + // Only a budget spent before the first match leaves nothing to serve. + // Discarding a prefix the walk already read strands the consumer: it + // polls the same offset again and is refused again. + if walk.limited && walk.matched == 0 { return DiskReadOutcome::Limited; } if walk.matched > 0 { @@ -772,6 +792,7 @@ impl DiskReadPlan { fragments: walk.fragments, last_matching_offset: walk.last_matching_offset, matched: walk.matched, + limited: walk.limited, faulted, } } else if faulted { @@ -809,7 +830,7 @@ impl DiskReadPlan { > walk.max_bytes { walk.limited = true; - return SegmentWalk::Faulted; + return SegmentWalk::Limited; } let Some(chunk) = self.read_chunk_with_retry(file, len, walk).await else { // Chunk read exhausted retries: same fail-closed reason as @@ -850,7 +871,7 @@ impl DiskReadPlan { } if walk.retained_bytes > walk.max_bytes { walk.limited = true; - return SegmentWalk::Faulted; + return SegmentWalk::Limited; } if corrupt { // A batch that does not match its own checksum. Fail closed like @@ -1742,13 +1763,77 @@ mod tests { assert!(matches!( plan.walk_segment(&file, offset_query(0), 1, bytes.len() as u64, &mut walk) .await, - SegmentWalk::Faulted + SegmentWalk::Limited )); assert!(walk.limited); assert_eq!(walk.matched, 0); assert!(walk.fragments.is_empty()); } + #[compio::test] + async fn a_byte_limited_disk_walk_serves_the_prefix_it_already_read() { + // The budget admits the first chunk and stops the next one. Before this + // was fixed the walk threw away every batch it had matched, so the + // consumer polled the same offset again and was refused again. + const BATCHES: u32 = 128; + let directory = tempfile::tempdir().unwrap(); + let mut bytes = Vec::new(); + for offset in 0..u64::from(BATCHES) { + bytes.extend(disk_batch(16 << 10, offset)); + } + std::fs::write(directory.path().join("00000000000000000000.log"), &bytes).unwrap(); + let plan = PollPlan { + context: PollContext { + history: PollHistoryId::default(), + consumer: PollingConsumer::Consumer(0, 0), + auto_commit: true, + }, + commit_offset: u64::from(BATCHES) - 1, + tier: PollTier::Disk { + query: MessageLookup::Offset { + offset: 0, + count: BATCHES, + ceiling: u64::from(BATCHES) - 1, + }, + resident_tail: ResidentTailSnapshot { + entries: Vec::new(), + oldest_resident: None, + }, + disk: DiskReadPlan { + partition_dir: PartitionDirResolution::Resolved( + directory.path().display().to_string(), + ), + bytes_per_message: None, + widest_batch_bytes: 0, + segments: vec![DiskSegment { + start_offset: 0, + persisted: bytes.len() as u64, + read_state: Rc::default(), + sealed: false, + }], + start_position: 0, + start_index_offset: None, + namespace_raw: 0, + validate_checksum: true, + }, + }, + }; + // Chunks are 1 MiB here, and the walk refuses a chunk whose buffer + // would put four times its length over the budget. This admits the + // first chunk and stops the second, which is the case under test. + let result = plan.execute_with_limit(4608 << 10).await; + let served = result.message_count(); + assert!( + served > 0 && served < BATCHES, + "expected a partial prefix, got {served} of {BATCHES}" + ); + assert_eq!(result.last_matching_offset, Some(u64::from(served) - 1)); + let result = result + .checked() + .expect("a byte-limited prefix is serveable"); + assert_eq!(result.message_count(), served); + } + fn resident(fragments: PollFragments) -> PollReadResult { PollPlan { context: PollContext {
