Skip block copies for spool stages whose workers are all remote - #19428
Conversation
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 Report✅ All modified and coverable lines are covered by tests. 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
Flags with carried forward coverage won't be shown. Click here to find out more. ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
gortiz
left a comment
There was a problem hiding this comment.
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
HashExchangepartitions rows disjointly, so its own workers never share a cell object — but the stage as a whole still keeps a reference alive pastsend(), which is what the outer exchange needs to know. Same forRandomExchange/SingletonExchange. Only an innerBroadcastExchangestrictly requires the OR; for the others it is stricter than necessary and never unsafe. GrpcSendingMailboxreally earns itsfalse.send→sendInternal→processAndSendserializes intoByteStrings 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:
-
The new fast path probably has no end-to-end coverage.
WindowFunnelTestandSpoolIntegrationTestboth 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 reportingdeliversByReference() == trueand 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. -
The
sendcontract is now load-bearing for query correctness and nothing enforces it. If an implementation returnsfalseand does not finish reading insidesend(), the result is silently wrong aggregates, not an exception.BlockExchangeSendingMailbox'sfalseis an inductive claim that rests onBlockExchange#routebeing 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 onroute, where it could actually be violated. -
Undocumented implication:
deliversByReference() == trueimpliesisLocal() == truefor 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. -
Minor: please put
#19427in theisLocal()TODO — it is in the description but not in the code. And a half-sentence on the cached_deliversByReferencenoting 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".
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.
AI-generated · Green: added · Yellow: modified · Red: removed · Gray: existing
Diff evidence
Fixes #19427. Follow-up to #19353.
Problem
#19353 made
BroadcastExchangecopy 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 withinsend.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, andisLocal()on that wrapper returnstrueunconditionally. 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-nullOBJECTcell.The cause is that
isLocal()answered two different questions:BlockExchange#sendBlockasks "can I pass the block whole, without splitting it?"BroadcastExchange#routeasks "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 inBroadcastExchange#route:InMemorySendingMailboxreturnstrue.GrpcSendingMailboxreturnsfalse.BlockExchangeSendingMailboxreturnstrueif any mailbox of its inner exchange returnstrue.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
HashExchangebuilds 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 whatsendBlockneeds. This PR also moves the contract thatroutedepends on — a mailbox that does not deliver by reference must finish reading the block beforesendreturns — fromGrpcSendingMailboxtoSendingMailbox#send, where implementers can see it.Testing
BroadcastExchangeTestgets 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.BlockExchangeTestcovers 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.
WindowFunnelTestandSpoolIntegrationTestrun 2 servers, where every stage normally has a worker on both, so their spool wrappers still reportdeliversByReference() == trueand 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() == truetoBlockExchange#sendBlock, while the inner exchanges getBlockSplitter.NO_OP, so multi-send blocks to remote receivers are never split againstMAX_MAILBOX_CONTENT_SIZE_BYTES. That defect is tracked separately in #19427.