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]

Reply via email to