Skip to content

[MSE] Add resumable sorted merge joins on shared exchanges - #19713

Open
xiangfu0 wants to merge 9 commits into
apache:xiangfu0/codex1/sorted-exchange-leaffrom
xiangfu0:xiangfu0/codex1/sorted-exchange-join
Open

xiangfu0 wants to merge 9 commits into
apache:xiangfu0/codex1/sorted-exchange-leaffrom
xiangfu0:xiangfu0/codex1/sorted-exchange-join

Conversation

@xiangfu0

@xiangfu0 xiangfu0 commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor

Adds opt-in resumable INNER and LEFT sorted merge joins. Explicit sorted strategy requires the broker capability gate, plans complete ASC NULLS LAST sender sorts and distinct k-way merge receives, and preserves explicit or inferred colocation through the Sort. Unsupported clusters fail logical planning; V2 keeps its explicit rejection.

Duplicate-key buffers and output blocks are bounded. Null/LONG/residual semantics remain intact. Early termination, exhausted inputs and BREAK drain terminal mailbox blocks so late errors and upstream statistics survive. The new leaf dependency's streaming instance-plan regression is preserved, and shared window fixtures have no incidental changes. The native diff contains 20 join feature files.

Validation: 716 unique applicable JDK 25 cases, including 35 sorted join tests, eight mailbox terminal regressions, three H2 comparisons, full following/RANGE window coverage and 57 core streaming combine cases; repeated classes are counted once and the two unchanged broker gate cases are explicitly reused. No failures/errors/skips or warnings on added lines. Mandatory affected formatting, license and checkstyle checks pass.

This PR remains stacked on #19712. Hosted CI at its final remote head and reviewer/base decisions remain open; local tests do not claim remote green CI. The cached capability gate is not an atomic downgrade guarantee, and #19395's stalled-sender memory followup remains deferred.

Final stacked head: b82686e30f759d1610e2021967c3cef7fd3234cd; dependency head: 88ef14a0fd97415f58c1aefc15c33a6620c759d6. Existing base names remain unchanged; the new dependency heads await root's guarded second publication.

The inherited fixture followup independently passes all 593 ResourceBasedQueryPlansTest cases on this exact source tree with compiler warnings/deprecations enabled, zero failures/errors/skips and all four planner formatting/license/style goals passing. Its native feature patch is byte-identical to the preceding published stack; only core owns the 124 corrected window-input fetch expectations. These local results do not establish hosted CI status.

The existing stack base is retained. Full unit/integration PR workflows are filtered to master; no complete hosted CI result exists for this child at its current stack base. After the prerequisite PRs land, the agreed base update and exact-head CI are still required. No merge, automatic merge, CI suppression or master integration was performed.

@xiangfu0
xiangfu0 force-pushed the xiangfu0/codex1/sorted-exchange-leaf branch from 42e3b46 to eecbbd1 Compare September 30, 2026 15:54
@xiangfu0
xiangfu0 force-pushed the xiangfu0/codex1/sorted-exchange-join branch from 33353c5 to ef0f689 Compare September 30, 2026 15:54
@xiangfu0
xiangfu0 marked this pull request as ready for review September 30, 2026 16:01
@xiangfu0
xiangfu0 force-pushed the xiangfu0/codex1/sorted-exchange-leaf branch from eecbbd1 to dbd31d9 Compare September 30, 2026 19:06
@xiangfu0
xiangfu0 force-pushed the xiangfu0/codex1/sorted-exchange-join branch 2 times, most recently from 1308a57 to f21d3c8 Compare October 3, 2026 19:45
@xiangfu0
xiangfu0 force-pushed the xiangfu0/codex1/sorted-exchange-leaf branch from dbd31d9 to 693f2f9 Compare October 3, 2026 19:45
rohityadav1993 and others added 8 commits October 4, 2026 04:20
Hash join materializes the entire right side into a hash table before probing,
so peak memory grows linearly with right-side row count, and multi-column keys
fall back to ObjectLookupTable with a composite key allocation per row. When
both inputs are already sorted on the join keys, a merge join avoids the hash
table entirely and streams.

- SortedMergeJoinOperator performs a two-pointer merge with lazy block reads
  (one block per side held in memory), type-dispatched key comparison, non-equi
  residual filter support, split equi-only/filter paths, maxRowsInJoin overflow
  handling (THROW/BREAK) on both the buffered right-key run and the emitted
  rows, periodic termination checks, and early-termination propagation so a
  downstream LIMIT stops the join instead of completing the full cross-product
  of buffered input.
- Selected via /*+ joinOptions(join_strategy='sorted') */, carried through the
  plan as JoinStrategy.SORTED.
- PinotJoinExchangeNodeInsertRule injects a LogicalSort below the sort exchange
  so both inputs arrive globally sorted, not merely sorted per sender. Join keys
  are collated NULLS LAST, matching the operator's key comparator.
- Colocation: joinOptions(is_colocated_by_join_keys='true') sets the existing
  prePartitioned field on the sort exchange, reusing the hash-join
  prePartitioned machinery, so co-partitioned inputs get a direct 1:1 exchange
  (receive fanIn drops to 1, no cross-server shuffle) with no changes needed in
  MailboxAssignmentVisitor or WorkerManager.

Part of apache#18667.

(cherry picked from commit 7590ae4)
Reproduce with a hot join key and LIMIT 1: the old merge built every match before returning. Resume in bounded blocks and limit buffered right rows instead of total output.
Cancel a duplicate-key join after its first output block to reproduce the retained input buffers. Drop the run and cursor rows on termination while preserving emitted rows and stats.
Add an H2 case using sortedSelectionMergeMode=on and streamingSortedMailboxReceive=true with a sorted LEFT join, residuals and ORDER BY LIMIT. Confirm the join ordering method implements the shared producer contract.
Materialize accepted equi and residual matches through the same row builder. Keep bounded INNER and LEFT joins, nulls, resource limits, cancellation and cleanup, with three H2 cases for joins, filters and ordered residual LEFT joins.
@xiangfu0
xiangfu0 force-pushed the xiangfu0/codex1/sorted-exchange-join branch from f21d3c8 to b82686e Compare October 3, 2026 20:25
@xiangfu0
xiangfu0 force-pushed the xiangfu0/codex1/sorted-exchange-leaf branch from 693f2f9 to 88ef14a Compare October 3, 2026 20:25
@codecov-commenter

codecov-commenter commented Oct 6, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 83.72881% with 48 lines in your changes missing coverage. Please review.
⚠️ Please upload report for BASE (xiangfu0/codex1/sorted-exchange-leaf@88ef14a). Learn more about missing BASE report.

Files with missing lines Patch % Lines
...uery/runtime/operator/SortedMergeJoinOperator.java 88.35% 18 Missing and 11 partials ⚠️
...not/query/runtime/operator/MultiStageOperator.java 14.28% 6 Missing ⚠️
...ite/rel/rules/PinotJoinExchangeNodeInsertRule.java 77.77% 0 Missing and 4 partials ⚠️
...e/pinot/query/runtime/InStageStatsTreeBuilder.java 20.00% 3 Missing and 1 partial ⚠️
.../apache/pinot/query/planner/plannode/JoinNode.java 40.00% 2 Missing and 1 partial ⚠️
.../query/planner/logical/RelToPlanNodeConverter.java 50.00% 0 Missing and 2 partials ⚠️
Additional details and impacted files
@@                           Coverage Diff                           @@
##             xiangfu0/codex1/sorted-exchange-leaf   #19713   +/-   ##
=======================================================================
  Coverage                                        ?   68.32%           
  Complexity                                      ?     1450           
=======================================================================
  Files                                           ?     3529           
  Lines                                           ?   230118           
  Branches                                        ?    36594           
=======================================================================
  Hits                                            ?   157226           
  Misses                                          ?    60559           
  Partials                                        ?    12333           
Flag Coverage Δ
integration 100.00% <ø> (?)
integration1 100.00% <ø> (?)
integration2 0.00% <ø> (?)
java-25 68.32% <83.72%> (?)
lane-a 100.00% <ø> (?)
lane-b 0.00% <ø> (?)
temurin 68.32% <83.72%> (?)
unittests 68.32% <83.72%> (?)
unittests1 58.41% <83.72%> (?)
unittests2 39.67% <1.35%> (?)

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.

@xiangfu0

xiangfu0 commented Oct 6, 2026

Copy link
Copy Markdown
Contributor Author

CI follow-up at c5852643 adds six exact stacked-base pull_request branch filters (+6/-0). Feature code/tests and existing master/push, path, permission and job rules are preserved.

Current-head checks finished with 13/14 passing. Integration Set 2 failed again after one failed-job-only rerun: SameTableNameMultiClusterIntegrationTest hit a broker HTTP BindException during setup, before query assertions (latest failed job). The first attempt failed at a different fixture listener (Netty QueryServer). All 11 directly implicated fixture/startup/POM files are identical to the frozen base; the actual conflicting socket remains unknown.

This PR still has an open CI gate. A separate baseline fixture startup repair and validation remain necessary; the six-line workflow change is preserved. Maintainer review remains pending.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants