Skip to content

Skip block copies for spool stages whose workers are all remote - #19428

Merged
yashmayya merged 2 commits into
apache:masterfrom
yashmayya:spool-delivers-by-reference
Sep 21, 2026
Merged

yashmayya merged 2 commits into
apache:masterfrom
yashmayya:spool-delivers-by-reference

Conversation

@yashmayya

@yashmayya yashmayya commented Sep 1, 2026 •

Copy link
Copy Markdown
Contributor

PR flow

BroadcastExchange now decides block copying based on whether any destination mailbox delivers blocks by reference, avoiding unnecessary copies for spool stages with all-remote workers.

flowchart TD
  N0["BroadcastExchange#35;route #40;F6#41;"]:::stModified
  N1["BlockExchangeSendingMailbox#35;deliversByReference #40;F5#41;"]:::stModified
  N2["BlockExchange#35;#95;deliversByReference #40;F5#41;"]:::stAdded
  N3["SendingMailbox#35;deliversByReference #40;interface#41; #40;F3#41;"]:::stAdded
  N4["GrpcSendingMailbox#35;deliversByReference #40;F1#41;"]:::stAdded
  N5["InMemorySendingMailbox#35;deliversByReference #40;F2#41;"]:::stAdded
  N0 -->|"calls"| N1
  N1 -->|"returns"| N2
  N2 -->|"anyDeliversByReference calls"| N3
  N3 -->|"implemented by"| N4
  N3 -->|"implemented by"| N5
  classDef stAdded fill:#dafbe1,stroke:#1a7f37,color:#1f2328,stroke-width:2px
  classDef stModified fill:#fff8c5,stroke:#9a6700,color:#1f2328,stroke-width:2px
  classDef stRemoved fill:#ffebe9,stroke:#cf222e,color:#1f2328,stroke-width:2px
  classDef stUnchanged fill:#f6f8fa,stroke:#656d76,color:#1f2328,stroke-width:1px
Loading

AI-generated · Green: added · Yellow: modified · Red: removed · Gray: existing

Diff evidence
  • F1: pinot-query-runtime/src/main/java/org/apache/pinot/query/mailbox/GrpcSendingMailbox.java — before · after
  • F2: pinot-query-runtime/src/main/java/org/apache/pinot/query/mailbox/InMemorySendingMailbox.java — before · after
  • F3: pinot-query-runtime/src/main/java/org/apache/pinot/query/mailbox/SendingMailbox.java — before · after
  • F5: pinot-query-runtime/src/main/java/org/apache/pinot/query/runtime/operator/exchange/BlockExchange.java — before · after
  • F6: pinot-query-runtime/src/main/java/org/apache/pinot/query/runtime/operator/exchange/BroadcastExchange.java — before · after
  • Regenerate PR flow

Fixes #19427. Follow-up to #19353.

Problem

#19353 made BroadcastExchange copy blocks that carry aggregation intermediate results, because a destination can hand an on-heap block to a receiver by reference. It skips the copy for destinations that are not local, because those serialize the block within send.

That skip never applied to spools, which are the case #19353 was written for. A multi-send (spool) node wraps each receiver stage's inner exchange as a BlockExchange.BlockExchangeSendingMailbox, and isLocal() on that wrapper returns true unconditionally. Every outer destination of a spool is such a wrapper, so a receiver stage with no worker on the sending server still got a copy. The inner exchange then serialized that copy immediately, on the same thread, and no receiver ever mutated it. Each wasted copy costs one serialization plus one deserialization of every non-null OBJECT cell.

The cause is that isLocal() answered two different questions:

  • BlockExchange#sendBlock asks "can I pass the block whole, without splitting it?"
  • BroadcastExchange#route asks "can a receiver keep a reference to this block?"

The two answers differ for BlockExchangeSendingMailbox. Only the conservative direction kept the code correct.

Fix

Add SendingMailbox#deliversByReference(), and use it in BroadcastExchange#route:

  • InMemorySendingMailbox returns true.
  • GrpcSendingMailbox returns false.
  • BlockExchangeSendingMailbox returns true if any mailbox of its inner exchange returns true.

The answer of a mailbox never changes, so each exchange computes it once, when it is created.

The delegation is an OR over the inner mailboxes. An inner HashExchange builds a new block for each destination, but those blocks hold the same cell objects, so one by-reference worker in a receiver stage is enough to require a copy.

isLocal() now only means what sendBlock needs. This PR also moves the contract that route depends on — a mailbox that does not deliver by reference must finish reading the block before send returns — from GrpcSendingMailbox to SendingMailbox#send, where implementers can see it.

Testing

BroadcastExchangeTest gets two spool cases: receiver stages whose workers are all remote (no copies), and a receiver stage with one worker on this server next to one on another server (one copy, shared within that stage). The first fails without this change. BlockExchangeTest covers the delegation directly.

The new fast path is covered by unit tests only. It fires when a receiver stage has no worker on the sending server. WindowFunnelTest and SpoolIntegrationTest run 2 servers, where every stage normally has a worker on both, so their spool wrappers still report deliversByReference() == true and still copy. Those suites show that this change causes no regression. They do not show that the fast path fires, and they stay green if it stops firing.

Out of scope

The same wrapper also reports isLocal() == true to BlockExchange#sendBlock, while the inner exchanges get BlockSplitter.NO_OP, so multi-send blocks to remote receivers are never split against MAX_MAILBOX_CONTENT_SIZE_BYTES. That defect is tracked separately in #19427.

BlockExchangeSendingMailbox is always local, but it only delivers by
reference when a mailbox of the exchange it decorates does. Reusing
isLocal() for the copy decision in BroadcastExchange therefore made every
spool receiver stage take a copy, including stages whose workers all
serialize the block.
@codecov-commenter

codecov-commenter commented Sep 1, 2026 •

Copy link
Copy Markdown

Codecov Report

✅ All modified and coverable lines are covered by tests.
✅ Project coverage is 67.85%. Comparing base (88ea694) to head (91ea27f).
⚠️ Report is 142 commits behind head on master.

Additional details and impacted files
@@             Coverage Diff              @@
##             master   #19428      +/-   ##
============================================
+ Coverage     67.57%   67.85%   +0.28%     
- Complexity     1430     1450      +20     
============================================
  Files          3487     3506      +19     
  Lines        224228   227060    +2832     
  Branches      35394    35894     +500     
============================================
+ Hits         151511   154063    +2552     
- Misses        60691    60863     +172     
- Partials      12026    12134     +108     
Flag Coverage Δ
integration 100.00% <ø> (ø)
integration1 100.00% <ø> (ø)
integration2 0.00% <ø> (ø)
java-25 67.85% <100.00%> (+0.28%) ⬆️
lane-a 100.00% <ø> (ø)
lane-b 0.00% <ø> (ø)
temurin 67.85% <100.00%> (+0.28%) ⬆️
unittests 67.84% <100.00%> (+0.28%) ⬆️
unittests1 58.02% <100.00%> (+0.32%) ⬆️
unittests2 39.55% <0.00%> (+0.23%) ⬆️

Flags with carried forward coverage won't be shown. Click here to find out more.

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@yashmayya yashmayya added enhancement Improvement to existing functionality multi-stage Related to the multi-stage query engine labels Sep 1, 2026
@yashmayya
yashmayya requested a review from gortiz September 1, 2026 18:41

@gortiz gortiz left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nice root-cause fix. Splitting isLocal() into the two questions it was conflating is the right call rather than special-casing the spool wrapper, and the docs on SendingMailbox are unusually clear about why the two-phase ordering in route needs both phases.

I walked the interleavings and found no correctness defect. Things I specifically checked:

  • The two-phase ordering survives the predicate swap. Copies are taken from the original inside the loop, at a point where phase 2 has not run, so no receiver can hold the original yet. I traced [mixed stage, all-remote stage, mixed stage]: the second mixed stage's inner exchange hands its copy to a local worker that may start mutating on the receiver thread while the outer loop is still running, but the all-remote stage has already finished reading the original and the first mixed stage has not received anything yet.
  • The OR delegation is conservative for every inner exchange type. An inner HashExchange partitions rows disjointly, so its own workers never share a cell object — but the stage as a whole still keeps a reference alive past send(), which is what the outer exchange needs to know. Same for RandomExchange/SingletonExchange. Only an inner BroadcastExchange strictly requires the OR; for the others it is stricter than necessary and never unsafe.
  • GrpcSendingMailbox really earns its false. send → sendInternal → processAndSend serializes into ByteStrings on the calling thread, and the back-pressure gate blocks before that rather than deferring the work, so the block is fully read on every path.

Non-blocking comments:

  1. The new fast path probably has no end-to-end coverage. WindowFunnelTest and SpoolIntegrationTest both run 2 servers, but the skip only fires when a receiver stage has no worker on the sending server — with 2 servers and broadcast/hash fan-out, every stage normally has a worker on both, so the spool wrappers keep reporting deliversByReference() == true and the copies still happen. Those tests prove "no regression"; they do not prove the optimization fires, and they would not catch it if the skip were wrong. Worth either asserting on a copy count, or just saying in the description that the fast path is unit-test-only, so a future refactor cannot silently reintroduce the copy with the suite green.

  2. The send contract is now load-bearing for query correctness and nothing enforces it. If an implementation returns false and does not finish reading inside send(), the result is silently wrong aggregates, not an exception. BlockExchangeSendingMailbox's false is an inductive claim that rests on BlockExchange#route being fully synchronous — true today for all four implementations, but it would break the moment any exchange queued a block for another thread. Consider restating it on route, where it could actually be violated.

  3. Undocumented implication: deliversByReference() == true implies isLocal() == true for all three implementations, and conceptually must (you cannot hand a heap object by reference across a socket). The reverse not holding is exactly the bug being fixed. Stating that in the javadoc would stop a future implementer producing the impossible remote-plus-by-reference combination.

  4. Minor: please put #19427 in the isLocal() TODO — it is in the description but not in the code. And a half-sentence on the cached _deliversByReference noting it deliberately ignores early termination (being wrong that way only costs a copy, never skips a needed one) would save the next reader the analysis.

Also worth a line in the comment on BlockExchangeSendingMailbox#deliversByReference: a reader who pictures HashExchange will think the OR is too strong, and the reason it is not is not obvious from "the decorated exchange passes blocks to its own mailboxes".

@yashmayya
yashmayya merged commit 12e7af1 into apache:master Sep 21, 2026
13 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

enhancement Improvement to existing functionality multi-stage Related to the multi-stage query engine

Projects

None yet

Development

Successfully merging this pull request may close these issues.

MSE: spool fan-out copies aggregation intermediates even when every receiver stage is remote

3 participants