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 3b726a65eb0bc2120b650dc1699e3eab14531ced Author: Hubert Gruszecki <[email protected]> AuthorDate: Thu Sep 17 18:51:51 2026 +0200 fix(shard): charge a deferred poll for its reply, not its borrows A deferred poll was refused with InvalidSizeBytes whenever its result borrowed more memory than the read had reserved. The charge counted a shared allocation once per fragment slicing it, and it ran before the reply was trimmed to the caller's byte limit, so a serveable poll failed on how the journal happened to allocate rather than on what it returned. Under load every poll on a busy partition failed, the client treated the refusal as a lost assignment, and the consumer fell behind for good. Count each backing allocation once, and take the charge after the trim. When the reply still pins more than the reservation, copy it into one exact-size allocation instead of refusing: fragment boundaries survive, so batch framing is unchanged, and the charge drops to the reply's own size. A poll that can return fewer messages must never fail. The check itself stays. It bounds real pinned memory, which the shard aggregates through max_inflight_bytes. --- core/partitions/src/poll_plan.rs | 101 ++++++++++++++++++++++++++++++++++++--- core/partitions/src/types.rs | 7 +++ core/shard/src/poll/deferred.rs | 14 +++--- 3 files changed, 109 insertions(+), 13 deletions(-) diff --git a/core/partitions/src/poll_plan.rs b/core/partitions/src/poll_plan.rs index 7306ea171..12d08342b 100644 --- a/core/partitions/src/poll_plan.rs +++ b/core/partitions/src/poll_plan.rs @@ -271,15 +271,50 @@ impl PollReadResult { Ok((self, limited)) } - /// Conservative charge: shared allocations count once for each fragment. + /// Memory this result keeps alive, counting every backing allocation once + /// however many fragments slice it. #[must_use] pub fn retained_bytes(&self) -> usize { - self.fragments.iter().fold( - self.fragments - .capacity() - .saturating_mul(size_of::<crate::Fragment>()), - |bytes, fragment| bytes.saturating_add(fragment.allocation_bytes()), - ) + let mut bytes = self + .fragments + .capacity() + .saturating_mul(size_of::<crate::Fragment>()); + for (index, fragment) in self.fragments.iter().enumerate() { + if !self.fragments[..index] + .iter() + .any(|earlier| earlier.shares_allocation_with(fragment)) + { + bytes = bytes.saturating_add(fragment.allocation_bytes()); + } + } + bytes + } + + /// Copy the selection into one exact-size allocation and drop the borrowed + /// buffers. A reply that slices a few bytes out of segment-sized journal + /// buffers pins all of them, so the copy trades one pass over the reply for + /// a charge that tracks the reply itself. Fragment boundaries survive, which + /// keeps batch framing intact. + #[must_use] + pub fn compacted(mut self) -> Self { + let total: usize = self.fragments.iter().map(Fragment::len).sum(); + let mut buffer = Owned::with_capacity(total); + for fragment in &self.fragments { + buffer.extend_from_slice(fragment.as_slice()); + } + let compacted = Frozen::from(buffer); + let mut start = 0; + self.fragments = self + .fragments + .iter() + .map(|fragment| { + let end = start + fragment.len(); + let piece = Fragment::slice(compacted.clone(), start, end); + start = end; + piece + }) + .collect(); + self } #[must_use] @@ -1720,4 +1755,56 @@ mod tests { assert_eq!(walk.matched, 0); assert!(walk.fragments.is_empty()); } + + fn resident(fragments: PollFragments) -> PollReadResult { + PollPlan { + context: PollContext { + history: PollHistoryId::default(), + consumer: PollingConsumer::Consumer(0, 0), + auto_commit: true, + }, + commit_offset: 0, + tier: PollTier::Resident { + fragments, + last_matching_offset: None, + message_count: 0, + }, + } + .execute_resident() + } + + fn windows(source: &Frozen<4096>, count: usize) -> PollFragments { + (0..count) + .map(|window| Fragment::slice(source.clone(), window * 4096, window * 4096 + 512)) + .collect() + } + + #[test] + fn retained_bytes_counts_one_allocation_once_however_many_slices_borrow_it() { + let source: Frozen<4096> = Owned::copy_from_slice(&vec![7; 1 << 20]).into(); + let one = resident(windows(&source, 1)).retained_bytes(); + let four = resident(windows(&source, 4)).retained_bytes(); + assert_eq!(one, four); + assert!(four < 2 * source.allocation_bytes()); + } + + #[test] + fn compacting_releases_the_borrowed_allocation_and_preserves_the_selection() { + let source: Frozen<4096> = Owned::copy_from_slice(&vec![7; 1 << 20]).into(); + let borrowed = resident(windows(&source, 4)); + let selection: Vec<Vec<u8>> = borrowed + .fragments + .iter() + .map(|fragment| fragment.as_slice().to_vec()) + .collect(); + let charge = borrowed.retained_bytes(); + let compacted = borrowed.compacted(); + assert!(compacted.retained_bytes() < charge); + let kept: Vec<Vec<u8>> = compacted + .fragments + .iter() + .map(|fragment| fragment.as_slice().to_vec()) + .collect(); + assert_eq!(kept, selection); + } } diff --git a/core/partitions/src/types.rs b/core/partitions/src/types.rs index 7619552c7..50a08570e 100644 --- a/core/partitions/src/types.rs +++ b/core/partitions/src/types.rs @@ -102,6 +102,13 @@ impl<const ALIGN: usize> Fragment<ALIGN> { pub fn borrows_from(&self, source: &Frozen<ALIGN>) -> bool { self.source.shares_allocation(source) } + + /// Whether both fragments keep the same allocation alive, so a byte charge + /// counts it once instead of once per slice. + #[must_use] + pub fn shares_allocation_with(&self, other: &Self) -> bool { + self.source.shares_allocation(&other.source) + } } /// Arguments for polling messages from a partition. diff --git a/core/shard/src/poll/deferred.rs b/core/shard/src/poll/deferred.rs index d9c51e75d..4f7e96d98 100644 --- a/core/shard/src/poll/deferred.rs +++ b/core/shard/src/poll/deferred.rs @@ -626,12 +626,7 @@ where return; } }; - let config = *self.deferred_polls.config.borrow(); - if result.retained_bytes() > config.max_read_bytes { - self.reject_deferred_poll(id, &waiter, IggyError::InvalidSizeBytes); - return; - } - let (result, byte_limited) = + let (mut result, byte_limited) = match result.limit_bytes(waiter.request.options.max_bytes as usize) { Ok(result) => result, Err(error) => { @@ -639,6 +634,13 @@ where return; } }; + // Charge what the reply carries, not what the read walked over. A + // selection that borrows journal buffers can still pin more than the + // reservation, so copy it out instead of refusing a serveable poll. + let config = *self.deferred_polls.config.borrow(); + if result.retained_bytes() > config.max_read_bytes { + result = result.compacted(); + } if result.message_count() >= waiter.request.options.min_count || byte_limited || waiter.terminal
