gortiz commented on PR #19396:
URL: https://github.com/apache/pinot/pull/19396#issuecomment-5951299327

   I'd like to use this PR to fix how MSE assigns the responsibility for 
ordering. Today, ordering should be a requirement that the logical plan 
guarantees. Instead, the physical plan and the runtime keep getting extra 
machinery to guarantee it: `sort`/`sortedOnSender` on the exchange and mailbox 
nodes, `SortedMailboxReceiveOperator`, and now a per-block `sorted.on.sender` 
marker. I consider the marker a blocker.
   
   The runtime has to trust the plan. Re-checking ordering on every block is 
like a JIT adding checks for invariants that the front-end compiler already 
proved. It is also fragile: `MailboxSendOperator.isSortedOnSender()` requires a 
`SortOperator` right below the send, so any order-preserving operator we add 
there later (a buffer, for example) silently disables the optimization. And it 
goes through the whole mailbox layer (`GrpcSendingMailbox`, 
`InMemorySendingMailbox`, `MailboxService`, `ReceivingMailbox`, 
`MailboxContentObserver`, `ChannelUtils`, `BlockingMultiStreamConsumer`) for an 
opt-in feature.
   
   What I'd suggest instead is to express ordering with simple plan nodes, not 
flags:
   
   - The sender has no ordering flags. If the receiver needs ordered streams, 
the plan guarantees the order upstream: usually with a Sort, but also by 
construction, for example after a sorted join.
   - The receiver uses a new plan node: a receive that k-way merges streams 
that are each already sorted, with an optional limit/offset. That's 
`SortedMailboxMergeReceiveOperator`, and the optional limit also covers the 
streaming ORDER BY LIMIT in #19711.
   - When the plan doesn't guarantee order upstream, it uses a plain receive 
and a Sort downstream, as #19412 already does.
   
   On the Calcite side, this could be a new `PinotKWayMergeSortExchange`: an 
exchange that requires its input to be sorted on the collation and preserves 
that order in its output. `PlanFragmenter` would turn it into a plain 
`MailboxSendNode` plus the new merge receive node. A rule similar to 
`PinotSortExchangeCopyRule` could push the limit/offset of a Sort above it into 
the receive.
   
   In practice, for this PR that means removing the mailbox metadata and the 
`sortedOnSender` parameter of `offer()`/`offerRaw()`, 
`MseBlockWithStats.isSortedOnSender()`, `isLastBlockSortedOnSender()`, 
`MailboxSendOperator.isSortedOnSender()`, `hasExplicitSortInput()` and the 
full-sort fallback in the merge receiver. A per-row order check is fine in 
tests, but not as a production path.
   
   Compatibility gets simpler too. On a new server:
   
   | Plan | Comes from | New server |
   |---|---|---|
   | plain receive | current broker | `MailboxReceiveOperator` |
   | receive with `sort=true` | broker older than #19412 | 
`MailboxReceiveOperator` + `SortOperator` |
   | new merge receive node | new broker, option enabled | k-way merge |
   
   For a new broker with old servers, it's the broker's responsibility not to 
send the new node and to keep generating a plan that old servers understand, 
i.e. a plain receive with a Sort downstream. We already have the mechanism for 
this: `SendStatsPredicate` in `SAFE` mode watches the instance configs and 
detects when any broker or server reports a different Pinot version. Today only 
servers use it to decide whether to send MSE stats, but the broker can reuse 
the same logic and only pick `PinotKWayMergeSortExchange` when the cluster is 
homogeneous. As a last line of defense, an old server can't deserialize the new 
node (`PlanNodeDeserializer` throws `Unsupported PlanNode type`), so the query 
fails loudly instead of returning wrong results.
   
   With this, `sortedOnSender` and `sort` in `plan.proto` and the flags in 
`PinotLogicalSortExchange` can be deprecated, and 
`SortedMailboxReceiveOperator` can go away. Old ORDER BY plans already have a 
Sort with the fetch above the receive, so they get the bounded implementation 
with a LIMIT. For old window plans, the server injects a Sort without a fetch. 
That part can be a follow-up.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to