Skip to content

[MSE] Stream ordered leaf selection across sorted segments - #19712

Open
xiangfu0 wants to merge 11 commits into
apache:xiangfu0/codex1/sorted-exchange-orderingfrom
xiangfu0:xiangfu0/codex1/sorted-exchange-leaf
Open

xiangfu0 wants to merge 11 commits into
apache:xiangfu0/codex1/sorted-exchange-orderingfrom
xiangfu0:xiangfu0/codex1/sorted-exchange-leaf

Conversation

@xiangfu0

@xiangfu0 xiangfu0 commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor

Streams selection ORDER BY over physically sorted segments and combines results incrementally. ON selects the streaming-capable path; AUTO resolves once from segment metadata before both operators are built; OFF preserves the existing path. Unsorted or incompatible segments retain the materialized fallback. Nulls, DESC, ties, pruning, cancellation/errors and segment ownership are preserved.

This PR builds on #19711's explicit sorted-input merge contract, including the release-version capability gate and finite LIMIT/OFFSET after merging. Its actual streaming instance-plan regression covers ON/AUTO, complete leaf ordering, the LIMIT+OFFSET budget and bounded output blocks. Shared window regressions match the dependency exactly. The native diff contains only 13 leaf feature files.

Validation: 625 affected JDK 25 warning/deprecation-enabled cases, 376 overlapping full planner/rule/window cases, and the shared-rule cleanup test; zero failures/errors/skips or added-line warnings. Applicable formatting, license and checkstyle checks pass. Final source tree matches the validated canonical input after test-only harmonization.

The PR remains stacked on the ordering alias. Hosted CI at its final remote head, reviewer decisions and documented AUTO reversibility/performance and segment-count followups remain open. The core stalled-sender memory issue stays deferred in #19395. No local result is described as remote green CI.

Final stacked head: 88ef14a0fd97415f58c1aefc15c33a6620c759d6; dependency head: 755721f7fd1548f3e94f81c86f4f64cfae4c66c8. 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-ordering branch from 5dc4fb3 to 697a710 Compare September 30, 2026 15:54
@xiangfu0
xiangfu0 force-pushed the xiangfu0/codex1/sorted-exchange-leaf branch from 42e3b46 to eecbbd1 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-ordering branch from 697a710 to 986291e Compare September 30, 2026 19:06
@xiangfu0
xiangfu0 force-pushed the xiangfu0/codex1/sorted-exchange-leaf branch 2 times, most recently from dbd31d9 to 693f2f9 Compare October 3, 2026 19:45
@xiangfu0
xiangfu0 force-pushed the xiangfu0/codex1/sorted-exchange-ordering branch from 986291e to 8e86cc7 Compare October 3, 2026 19:45
rohityadav1993 and others added 10 commits October 4, 2026 04:20
An unbounded leaf-stage ORDER BY (as injected for sorted merge join inputs)
routes to MinMaxValueBasedSelectionOrderByCombineOperator, which merges every
segment's rows into a single block before returning anything. At large data
volumes this exceeds the leaf stage's CPU budget and ThreadAccountant raises
EarlyTerminationException inside SelectionOperatorUtils.mergeWithOrdering(),
surfacing to the broker as a spurious "Cancelled by sender".

This adds a streaming alternative, opt-in via the `streamingSelectionOrderBy`
query option:

- StreamingSelectionOrderByOperator emits sorted blocks incrementally for a
  segment that is physically sorted on the leading ORDER BY column, reading
  the sorted forward index in order instead of building a priority queue.
- StreamingSelectionOrderByCombineOperator performs a k-way heap merge across
  segment operators and emits bounded blocks (`streamingSelectionOrderByBlockSize`,
  default 10000) rather than one materialized result.
- SelectionPlanNode and CombinePlanNode select these operators when the option
  is set and the sortedness precondition holds; otherwise behaviour is unchanged.

Part of apache#18667.
Eight rounds of review on apache#19120, folded into one commit.
The streaming merge stays opt-in and off by default, so a query that
does not ask for it is unaffected by everything below.

- Restrict the path to plans that have a ResultsBlockStreamer. On the
  blocking path the caller takes a single nextBlock() and expects the
  whole result, so the merge could never flush and would hold every
  activated cursor's block alive for nothing.
- Replace the boolean streamingSelectionOrderBy option with the
  tri-valued sortedSelectionMergeMode (OFF / ON / AUTO, default OFF;
  the block size option is renamed to match), and resolve AUTO once in
  InstancePlanMakerImplV2 rather than at the gates. CombinePlanNode
  builds its leaves before evaluating its own gate, so a decision there
  arrives too late to stop streaming leaves being built under a
  non-streaming combine. AUTO takes the merge only when at least
  sortedSelectionMergeAutoMinSortedRatio of the queried segments are
  physically sorted on the leading ORDER BY column, decided from
  segment metadata alone before any segment is acquired.
- Resolve AUTO to OFF for a DESC leading expression without
  allowReverseOrder. DataSourceMetadata.isSorted() reports ascending
  physical order, so such a query cannot stream, and taking the path
  would give up the MinMax combine's parallelism and min/max pruning
  in exchange for nothing.
- Route both scan paths through one empty-block skip. The two paths
  disagreed about the project operator's contract: the sorted path
  read an empty block as exhaustion and would have ended a scan
  mid-segment.
- Emit one schema-carrying empty block from a segment that matches no
  rows, so a consumer can always learn the child's DataSchema, and
  derive the two-phase schema from column metadata instead of standing
  up a projection and transform operator per empty segment.
- Hold the winning cursor outside the heap. Segments are near-disjoint
  on a time-like leading column, so one cursor supplies long runs of
  consecutive rows, each of which previously paid two O(log k) sift
  operations to arrive back at the same cursor.
- Return a growable row list from the tail-to-sort drain, and fail
  loud on the two base worker entry points the combine replaces, so a
  future change routing back through them breaks at the seam instead
  of dereferencing the null merger at query time.
- Reuse the sorted-by project on the DESC fall-through instead of
  building a second identical one.

Every behavioural change above carries a regression test verified to
fail against the unfixed code.
Surfaced by the arm-4 memory benchmark, not by review. Where segment
min/max values tie on the leading ORDER BY column -- the normal shape
for a low-cardinality or timestamp-prefix sort key -- every cursor
activated at once, each pinning a decompressed block, because
sortsBeyond() only deferred a cursor sorting strictly past the merge
frontier. MinMaxValueBasedSelectionOrderByCombineOperator has had the
tie case since it was written; this path did not.

Relax the comparison to non-strict, gated on a single order-by
expression. With two or more, a column-0 tie can hide a row sorting
earlier on column 1, which would be genuine out-of-order emission --
the same gate the MinMax operator applies for the same reason.

The justification differs from that operator's, though, and the javadoc
says so: its bound is the k-th row of a complete top-K, so a tie
provably cannot improve the answer. Here the bound is the live merge
frontier, taken while fewer than limit + offset rows have been emitted,
so that argument is unavailable. What holds instead is that this only
ever defers: _nextToActivate does not advance on a deferral, and the
caller force-activates once the heap and leader both drain. Since the
bound bounds every row in the segment, a deferred cursor holds no row
sorting strictly before one already emitted -- only rows tying it,
interchangeable when column 0 is the whole sort key.

The frontier is the smaller of the retained leader's head and the heap
top's head, since that is the next row emitted. Defer when the bound
sorts past either head, which is exactly past the smaller one; checking
against the larger overstated the frontier whenever the leader had moved
past the heap head, and every remaining tied cursor activated. A present
candidate with a null head still forces activation.

This does change which tied rows a single-column ORDER BY returns. The
class already declares tie order arbitrary for rows equal on every
order-by expression, which under one column is the same set.

Tests on the tied-minima fixture: the deferral itself ASC and DESC
(both fail without the change, scanning every segment instead of one),
value-multiset correctness across a LIMIT straddling a tied group and
across an OFFSET, the two-expression gate holding full row parity, no
under-delivery when the LIMIT covers every row, and one segment
activated per tied run once the leader moves past the heap head, ASC,
DESC, and across output blocks.
Test segment release and schema-mismatch handling. Nothing exercised either
path before: the existing tests plan the combine from plain segment plan
nodes, so the acquire/release calls were no-ops and a leaked segment would
have passed, and every segment shared one schema, so the mismatch branch
never ran. The new tests wrap each real streaming leaf in a counting
AcquireReleaseColumnsSegmentOperator, the way the prefetch path plans it,
and assert every acquired segment is released exactly once when the merge
drains, when the limit stops it with cursors still open, on stop()
mid-merge (including a repeated stop()), and when a child throws an
exception or an Error. A mismatched schema on one block must drop the rest
of that segment, leave every other segment whole, and report one
deduplicated MERGE_RESPONSE error. Each was checked against an operator
with the corresponding release or drop removed. Document why the cursor
drops the rest of the segment on a mismatch rather than one block: it is
the same unit the merger drops, since there a block is a segment's whole
result.

Add server config defaults for the merge options.
sortedSelectionMergeBlockSize and sortedSelectionMergeAutoMinSortedRatio
could only be changed per query. Add
pinot.server.query.executor.sorted.selection.merge.block.size and
pinot.server.query.executor.sorted.selection.merge.auto.min.sorted.ratio,
read and validated in InstancePlanMakerImplV2#init and applied whenever the
query does not set the option, as numGroupsLimit already is. The block size
default moves from Broker to Server, since only the server reads it.
applyQueryOptions now always writes both values, so the AUTO threshold
tests set the ratio through the query option instead of the query context
setter. Also bound CombineSlowOperatorsTest's deadline test with a timeout:
a regression there drives a SlowOperator that sleeps for an hour, which
would hang the build instead of failing it.

Expose the combine's settings in the explain plan. It now reports block
size, frontier pruning, tie deferral and the segment / sorted-segment
counts as explain attributes, so an MSE explain with explainAskingServers
shows which path ran and whether pruning was active. Tests pin that
numSegmentsMatched excludes never-activated segments.

Use markdown syntax in the /// javadoc added by this change: [Foo],
backticks, **bold** and blank-line paragraphs instead of {@link}, {@code},
<b> and <p>, matching the rest of the module.
A query-level null-handling switch turned pruning off for every segment. The non-null metadata flag is already loaded with the segment, so flagged columns keep their min/max and unflagged ones always activate. AUTO uses the same predicate and no longer selects the streaming merge over segments the leaf will not stream. A consuming segment has no column metadata map, so it counts as unflagged rather than failing the query.
MinMax prunes on the metadata max, which ignores nulls, so with fewer worker threads than segments it can skip the
null-bearing segment and drop the null rows that DESC puts first. That made the parity check fail on CI's smaller
runners. Assert the expected rows directly.
Use one project-block and schema setup path, releasing no-tail block references after materialization. Keep fallback, ties, nulls and DESC coverage while replacing repeated empty-block fixtures with a small control over real segment reads.
@xiangfu0
xiangfu0 force-pushed the xiangfu0/codex1/sorted-exchange-leaf branch from 693f2f9 to 88ef14a Compare October 3, 2026 20:25
@xiangfu0
xiangfu0 force-pushed the xiangfu0/codex1/sorted-exchange-ordering branch from 8e86cc7 to 755721f Compare October 3, 2026 20:25
@codecov-commenter

codecov-commenter commented Oct 6, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 94.58545% with 32 lines in your changes missing coverage. Please review.
⚠️ Please upload report for BASE (xiangfu0/codex1/sorted-exchange-ordering@755721f). Learn more about missing BASE report.

Files with missing lines Patch % Lines
...rator/query/StreamingSelectionOrderByOperator.java 91.62% 11 Missing and 8 partials ⚠️
...bine/StreamingSelectionOrderByCombineOperator.java 96.71% 1 Missing and 6 partials ⚠️
...va/org/apache/pinot/core/plan/CombinePlanNode.java 33.33% 0 Missing and 2 partials ⚠️
.../org/apache/pinot/core/plan/SelectionPlanNode.java 94.11% 0 Missing and 2 partials ⚠️
...pinot/core/plan/maker/InstancePlanMakerImplV2.java 97.01% 1 Missing and 1 partial ⚠️
Additional details and impacted files
@@                             Coverage Diff                             @@
##             xiangfu0/codex1/sorted-exchange-ordering   #19712   +/-   ##
===========================================================================
  Coverage                                            ?   68.30%           
  Complexity                                          ?     1450           
===========================================================================
  Files                                               ?     3528           
  Lines                                               ?   229831           
  Branches                                            ?    36530           
===========================================================================
  Hits                                                ?   156984           
  Misses                                              ?    60533           
  Partials                                            ?    12314           
Flag Coverage Δ
integration 100.00% <ø> (?)
integration1 100.00% <ø> (?)
integration2 0.00% <ø> (?)
java-25 68.30% <94.58%> (?)
lane-a 100.00% <ø> (?)
lane-b 0.00% <ø> (?)
temurin 68.30% <94.58%> (?)
unittests 68.30% <94.58%> (?)
unittests1 58.36% <94.58%> (?)
unittests2 39.72% <3.04%> (?)

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 83040faa 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.

All 14 current-head checks passed: unit/integration, compatibility, quickstart, Java 11 client, linter and Trivy (checks). Maintainer review remains pending.

The first attempt hit an unchanged ingestion-test Mockito concurrency error; one failed-job-only rerun passed on the same head (original attempt).

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