This is an automated email from the ASF dual-hosted git repository. xiangfu0 pushed a commit to branch xiangfu0/codex1/sorted-exchange-leaf in repository https://gitbox.apache.org/repos/asf/pinot.git
commit 0da93353bfcf305806218c14e6bb241ea52fcb90 Author: rohity <[email protected]> AuthorDate: Sat Sep 26 18:38:51 2026 +0000 Defer tied cursors under a single-column ORDER BY 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. --- .../StreamingSelectionOrderByCombineOperator.java | 50 ++++-- ...reamingSelectionOrderByCombineOperatorTest.java | 171 +++++++++++++++++++++ 2 files changed, 208 insertions(+), 13 deletions(-) diff --git a/pinot-core/src/main/java/org/apache/pinot/core/operator/combine/StreamingSelectionOrderByCombineOperator.java b/pinot-core/src/main/java/org/apache/pinot/core/operator/combine/StreamingSelectionOrderByCombineOperator.java index 527e7a77fb7..75560d3715b 100644 --- a/pinot-core/src/main/java/org/apache/pinot/core/operator/combine/StreamingSelectionOrderByCombineOperator.java +++ b/pinot-core/src/main/java/org/apache/pinot/core/operator/combine/StreamingSelectionOrderByCombineOperator.java @@ -110,6 +110,9 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi private final boolean _asc; private final boolean _pruningEnabled; + /// Whether a cursor whose bound *ties* the merge frontier may be deferred, not only one sorting strictly + /// beyond it. Only correct under a single order-by expression: see {@link #sortsBeyond}. + private final boolean _deferTiedCursors; private final int _numRowsToKeep; private final int _blockSize; private final Comparator<Object[]> _comparator; @@ -156,6 +159,10 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi firstOrderByExpressionContext.getType() == ExpressionContext.Type.IDENTIFIER ? firstOrderByExpressionContext.getIdentifier() : null; _pruningEnabled = firstOrderByColumn != null && !queryContext.isNullHandlingEnabled(); + // A column-0 tie only proves the segment cannot supply an *earlier* row when column 0 is the whole sort key. + // With two or more expressions a tie on column 0 can hide a row sorting earlier on column 1, so ties must + // still activate (same gate as MinMaxValueBasedSelectionOrderByCombineOperator's numOrderByExpressions == 1). + _deferTiedCursors = orderByExpressions.size() == 1; // Build one cursor per segment operator and read its first order-by column min/max for lazy activation ordering. // Reading DataSourceMetadata does not touch column buffers, so no segment acquire is needed here (mirrors @@ -325,11 +332,15 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi /// Cursors are visited in min/max order and {@code _nextToActivate} advances only on a real activation, so a /// {@code break} defers the current cursor to a later call with a risen frontier rather than skipping it. Deferring /// is safe because the first order-by column is the primary sort key: a cursor whose range starts past the frontier - /// holds no row sorting before it. With no frontier known yet, activation is forced. + /// holds no row sorting before it. With no frontier known yet, activation is forced. Under a single order-by + /// expression a cursor whose bound merely *ties* the frontier is deferred too, which is what keeps a long run of + /// equal segment minima from activating every segment at once; see {@link #sortsBeyond}. /// /// The frontier is the row about to be emitted. Understating it defers a cursor that could have supplied that row /// and the merge emits out of order, where overstating it only activates a segment early. The leader is retained - /// outside the heap, so that row is either its head or the heap head: a bound past both is past the frontier. + /// outside the heap, so that row is the smaller of its head and the heap head: a bound past either is past the + /// frontier. Checking the leader alone would overstate it whenever the leader has moved past the heap head, which + /// happens both on entry (the leader just advanced) and within this loop (an activation offers a smaller head). private void activateEligibleCursors(@Nullable SegmentCursor leader) { while (_nextToActivate < _sortedCursors.length) { SegmentCursor cursor = _sortedCursors[_nextToActivate]; @@ -338,7 +349,7 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi // Re-read per iteration: each activation below can offer a smaller head into the heap. SegmentCursor top = _priorityQueue.peek(); // A null bound always activates, and with no candidate at all there is no frontier to prune against. - if (bound != null && (leader != null || top != null) && sortsBeyond(bound, leader) && sortsBeyond(bound, top)) { + if (bound != null && (leader != null || top != null) && sortsBeyond(bound, leader, top)) { break; } } @@ -350,21 +361,34 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi } } - /// Whether `bound` sorts past the cursor's head on the first order-by column, so that segment cannot supply the - /// head. An absent cursor imposes no constraint; a null head cannot be compared, so this reports `false` and the - /// caller activates. Tests column 0 only, as the pruning bound always has, rather than the full-row comparator -- - /// which would put every order-by column on the per-row path. - private boolean sortsBeyond(Comparable bound, @Nullable SegmentCursor cursor) { - if (cursor == null) { - return true; - } - Object headValue = cursor.currentHead()[0]; - if (headValue == null) { + /// Whether `bound` sorts past the frontier -- the smaller of the two candidates' heads -- on the first order-by + /// column, so that segment cannot supply the next row. An absent candidate imposes no constraint; a null head cannot + /// be compared, so this reports `false` and the caller activates. Tests column 0 only, as the pruning bound always + /// has, rather than the full-row comparator -- which would put every order-by column on the per-row path. + /// + /// Under a single order-by expression ({@code _deferTiedCursors}) a *tie* also defers, which only postpones a + /// cursor: {@code _nextToActivate} does not advance, and the caller force-activates once heap and leader drain. + /// Since {@code bound} bounds every row in the segment, a deferred cursor holds no row sorting strictly before an + /// emitted one -- only ties, interchangeable when column 0 is the whole sort key. Note this is *not* + /// {@link MinMaxValueBasedSelectionOrderByCombineOperator}'s justification: its bound is a complete top-K's k-th + /// row, this one is the live frontier. It changes which tied rows are returned, already documented as arbitrary. + private boolean sortsBeyond(Comparable bound, @Nullable SegmentCursor leader, @Nullable SegmentCursor top) { + Object leaderHead = leader != null ? leader.currentHead()[0] : null; + Object topHead = top != null ? top.currentHead()[0] : null; + if ((leader != null && leaderHead == null) || (top != null && topHead == null)) { return false; } + // Past either head is past the smaller one. + return (leader != null && sortsBeyond(bound, leaderHead)) || (top != null && sortsBeyond(bound, topHead)); + } + + private boolean sortsBeyond(Comparable bound, Object headValue) { // Both come from the same first order-by column: the metadata min/max and the materialized row[0] share the // column's stored type, so this comparison is type-safe (same assumption as MinMaxValueBased...). int cmp = bound.compareTo(headValue); + if (_deferTiedCursors) { + return _asc ? cmp >= 0 : cmp <= 0; + } return _asc ? cmp > 0 : cmp < 0; } diff --git a/pinot-core/src/test/java/org/apache/pinot/core/operator/combine/StreamingSelectionOrderByCombineOperatorTest.java b/pinot-core/src/test/java/org/apache/pinot/core/operator/combine/StreamingSelectionOrderByCombineOperatorTest.java index 0ba5d5eb5bf..ae269cc1e76 100644 --- a/pinot-core/src/test/java/org/apache/pinot/core/operator/combine/StreamingSelectionOrderByCombineOperatorTest.java +++ b/pinot-core/src/test/java/org/apache/pinot/core/operator/combine/StreamingSelectionOrderByCombineOperatorTest.java @@ -390,6 +390,167 @@ public class StreamingSelectionOrderByCombineOperatorTest { "SELECT sortedCol, tailCol, valCol FROM testTable ORDER BY sortedCol, valCol LIMIT 40", false); } + /// Pins the tie-deferral itself: under a single order-by expression, every `_lowCardSegments` segment's sortedCol + /// minimum ties at 0, so a LIMIT satisfiable from one segment's leading run of 0s (25 rows) must not activate any + /// other segment. Before `_deferTiedCursors` existed, every cursor tied the (unknown-yet) frontier at activation + /// time and all of them activated up front, each scanning its first block (about LIMIT docs), so docs scanned + /// pre-fix is about `NUM_SEGMENTS * 20`. Post-fix only the leading segment (in cursor order) is ever touched, so docs + /// scanned is bounded by one segment's worth. + @Test + public void testTieDeferralBoundsDocsScannedToLeadingSegment() { + Result result = run(_lowCardSegments, "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol LIMIT 20", true, + false, true, 0); + assertTrue(result._combineOperator instanceof StreamingSelectionOrderByCombineOperator); + assertEquals(result._rows.size(), 20); + for (Object[] row : result._rows) { + assertEquals((int) row[0], 0, "LIMIT 20 must be satisfiable entirely from sortedCol=0 rows"); + } + assertTrue(result._numDocsScanned <= NUM_RECORDS_PER_SEGMENT, + "A LIMIT satisfiable from one segment's leading run must not activate any other tied segment; docs scanned: " + + result._numDocsScanned); + } + + /// DESC counterpart of [#testTieDeferralBoundsDocsScannedToLeadingSegment]: every `_lowCardSegments` segment's + /// sortedCol maximum ties at 3, so `sortsBeyond`'s `_asc ? cmp >= 0 : cmp <= 0` DESC branch is the one under test + /// here rather than the ASC one above. A sign flip in that branch would either over-scan (fail this test the same + /// way the pre-fix code fails the ASC test) or, worse, prune a segment that still had rows to give -- which the + /// row-count / row-value assertions below would catch. + @Test + public void testTieDeferralBoundsDocsScannedToLeadingSegmentDesc() { + Result result = run(_lowCardSegments, + "SET allowReverseOrder=true; SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol DESC LIMIT 20", true, + false, true, 0); + assertTrue(result._combineOperator instanceof StreamingSelectionOrderByCombineOperator); + assertEquals(result._rows.size(), 20); + for (Object[] row : result._rows) { + assertEquals((int) row[0], 3, "LIMIT 20 DESC must be satisfiable entirely from sortedCol=3 rows"); + } + assertTrue(result._numDocsScanned <= NUM_RECORDS_PER_SEGMENT, + "A LIMIT satisfiable from one segment's leading run must not activate any other tied segment; docs scanned: " + + result._numDocsScanned); + } + + /// Correctness under single-column ties, LIMIT straddling a tied group: with `ORDER BY sortedCol` alone every + /// segment ties at every one of the 4 distinct values, so a LIMIT of 30 forces the merge past the sortedCol=0 + /// boundary (only 25 such rows per segment) partway through a block from the leading segment. A wrongly-timed + /// reactivation here would either drop true sortedCol=0 rows from other segments (undercount at the tie) or emit + /// out of order. [#assertParity]'s full-row multiset check does not apply: with sortedCol as the only key, a + /// different (equally valid) tied row may be returned than the MinMax baseline picks, so only the order-by column + /// values -- whose multiset *is* pinned by the LIMIT boundary, unlike the underlying rows -- are compared. + @Test + public void testSingleColumnTieDeferralPreservesOrderByValuesAcrossLimitStraddle() { + @Language("sql") String query = "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol LIMIT 30"; + Result baseline = run(_lowCardSegments, query, false, false, false, 0); + Result streamed = run(_lowCardSegments, query, true, false, true, 3); + assertTrue(streamed._combineOperator instanceof StreamingSelectionOrderByCombineOperator); + assertEquals(streamed._rows.size(), 30); + assertSorted(streamed._rows, orderByComparator(query, false)); + assertEquals(orderByColumnValues(streamed._rows), orderByColumnValues(baseline._rows), + "Multiset of sortedCol values must match the MinMax baseline even though the underlying rows may differ"); + } + + /// Correctness under single-column ties with an OFFSET: LIMIT 10 OFFSET 20 straddles the same sortedCol=0/1 + /// boundary as [#testSingleColumnTieDeferralPreservesOrderByValuesAcrossLimitStraddle] (limit + offset = 30), but + /// exercises it with a nonzero offset. As [#testLimitOffsetParity] documents, the server (and this combine) retains + /// {@code limit + offset} rows -- the offset is trimmed by the broker afterwards -- so 30, not 10, rows come back + /// here too; what differs from the straddle test is only the query shape, not the row count. + @Test + public void testSingleColumnTieDeferralPreservesOrderByValuesAcrossOffset() { + @Language("sql") String query = "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol LIMIT 10 OFFSET 20"; + Result baseline = run(_lowCardSegments, query, false, false, false, 0); + Result streamed = run(_lowCardSegments, query, true, false, true, 3); + assertTrue(streamed._combineOperator instanceof StreamingSelectionOrderByCombineOperator); + assertEquals(streamed._rows.size(), 30, "Server retains limit + offset rows; the broker trims the offset later"); + assertSorted(streamed._rows, orderByComparator(query, false)); + assertEquals(orderByColumnValues(streamed._rows), orderByColumnValues(baseline._rows), + "Multiset of sortedCol values must match the MinMax baseline even though the underlying rows may differ"); + } + + /// Condition 1 (the two-expression gate): `_lowCardSegments` still tie on sortedCol alone, but a second order-by + /// expression (valCol) makes the full order-by key a total order, so full-row parity applies unlike the + /// single-column tests above. `_deferTiedCursors` must be false here (`orderByExpressions.size() == 1` fails), so + /// this pins the same shape both before and after the production change: a wrongly-deferred cursor on this shape + /// would drop a row that sorts earlier on valCol despite tying on sortedCol, which parity would catch as a missing + /// row rather than merely a reordered one. + @Test + public void testTwoExpressionOrderByDoesNotDeferTiedCursors() { + assertParity(_lowCardSegments, "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol LIMIT 30", + false); + } + + /// Condition 2 (no under-delivery when the heap drains): a single-column ORDER BY with a LIMIT covering every row + /// of `_lowCardSegments` forces every segment to eventually activate no matter how aggressively ties are deferred + /// -- deferral only postpones a cursor, it never removes it from consideration, and the + /// `(leader != null || top != null)` guard in `activateEligibleCursors` forces activation once the heap and leader + /// both go empty. Asserts the full row count comes back with none missing. + @Test + public void testSingleColumnTieDeferralNeverUnderDeliversAtFullDrain() { + int totalRows = NUM_SEGMENTS * NUM_RECORDS_PER_SEGMENT; + @Language("sql") String query = + "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol LIMIT " + totalRows; + Result baseline = run(_lowCardSegments, query, false, false, false, 0); + Result streamed = run(_lowCardSegments, query, true, false, true, 7); + assertTrue(streamed._combineOperator instanceof StreamingSelectionOrderByCombineOperator); + assertEquals(streamed._rows.size(), totalRows, "Every row must be delivered when the LIMIT covers the whole set"); + assertSorted(streamed._rows, orderByComparator(query, false)); + assertEquals(orderByColumnValues(streamed._rows), orderByColumnValues(baseline._rows), + "Multiset of sortedCol values must match the MinMax baseline when the merge drains completely"); + } + + /// Each `_lowCardSegments` segment holds only 25 rows at the tied minimum, so a LIMIT past 25 moves the leader's head + /// to sortedCol=1 while the other segments still tie at 0. The next row out is then the smaller of the leader's and + /// the heap's heads, and a waiting cursor must be deferred against that row, not against the leader alone: checking + /// the leader activates every remaining tied segment the moment one of them has been opened. Activating one segment + /// per 25-row run is the minimum, so at most `ceil(limit / 25)` segments are matched (a segment is matched once + /// activation reads its first block). + @Test + public void testTieDeferralActivatesOneSegmentPerTiedRunPastLeaderRun() { + for (int limit : new int[]{30, 60, 90}) { + Result result = run(_lowCardSegments, "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol LIMIT " + limit, + true, false, true, 0); + assertTrue(result._combineOperator instanceof StreamingSelectionOrderByCombineOperator); + assertEquals(result._rows.size(), limit); + int numZeros = Math.min(limit, 25 * NUM_SEGMENTS); + assertEquals(result._rows.stream().filter(row -> (int) row[0] == 0).count(), numZeros, + "Every returned row must come from the tied minimum while enough such rows exist"); + int maxSegments = (limit + 24) / 25; + assertTrue(result._numSegmentsMatched <= maxSegments, + "LIMIT " + limit + " should activate at most " + maxSegments + " segments; activated: " + + result._numSegmentsMatched); + } + } + + /// DESC counterpart of [#testTieDeferralActivatesOneSegmentPerTiedRunPastLeaderRun]: the leader leaves the tied + /// maximum (3) for 2 while the waiting segments still tie at 3. + @Test + public void testTieDeferralActivatesOneSegmentPerTiedRunPastLeaderRunDesc() { + Result result = run(_lowCardSegments, + "SET allowReverseOrder=true; SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol DESC LIMIT 30", true, + false, true, 0); + assertTrue(result._combineOperator instanceof StreamingSelectionOrderByCombineOperator); + assertEquals(result._rows.size(), 30); + for (Object[] row : result._rows) { + assertEquals((int) row[0], 3, "LIMIT 30 DESC must be satisfiable entirely from sortedCol=3 rows"); + } + assertTrue(result._numSegmentsMatched <= 2, + "LIMIT 30 DESC should activate at most 2 segments; activated: " + result._numSegmentsMatched); + } + + /// Same shape with a 7-row output block, so the leader is handed back to the heap and re-chosen mid-run: deferral + /// against the next row must hold across block boundaries too, and the rows must still match the baseline's values. + @Test + public void testTieDeferralPastLeaderRunAcrossOutputBlocks() { + @Language("sql") String query = "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol LIMIT 60"; + Result baseline = run(_lowCardSegments, query, false, false, false, 0); + Result streamed = run(_lowCardSegments, query, true, false, true, 7); + assertTrue(streamed._combineOperator instanceof StreamingSelectionOrderByCombineOperator); + assertEquals(streamed._rows.size(), 60); + assertSorted(streamed._rows, orderByComparator(query, false)); + assertEquals(orderByColumnValues(streamed._rows), orderByColumnValues(baseline._rows)); + assertTrue(streamed._numSegmentsMatched <= 3, + "LIMIT 60 should activate at most 3 segments; activated: " + streamed._numSegmentsMatched); + } + @Test public void testTwoPhaseSelectNonOrderByParity() { // tailCol is selected but not an order-by key -> the streaming children take the two-phase (order-by-then-fetch) @@ -964,6 +1125,7 @@ public class StreamingSelectionOrderByCombineOperatorTest { result._schema = block.getDataSchema(); } result._numDocsScanned = block.getNumDocsScanned(); + result._numSegmentsMatched = block.getNumSegmentsMatched(); break; } SelectionResultsBlock dataBlock = (SelectionResultsBlock) block; @@ -1026,6 +1188,14 @@ public class StreamingSelectionOrderByCombineOperatorTest { }).sorted().collect(Collectors.toList()); } + /// Extracts the ORDER BY column (projected at index 0 in every query these tests use) as a sorted multiset, + /// ignoring every other projected column. Used where a full-row [#assertMultisetEquals] would be too strong: under + /// a single order-by expression a tied value may be satisfied by different underlying rows than the MinMax + /// baseline picks, but the count of each order-by value in the result is still pinned by the LIMIT boundary. + private static List<Integer> orderByColumnValues(List<Object[]> rows) { + return rows.stream().map(row -> (Integer) row[0]).sorted().collect(Collectors.toList()); + } + @AfterClass public void tearDown() throws IOException { @@ -1047,5 +1217,6 @@ public class StreamingSelectionOrderByCombineOperatorTest { private List<Integer> _blockSizes; private int _numBlocks; private long _numDocsScanned; + private int _numSegmentsMatched; } } --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
