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 eecbbd1c04489681085b64828c0118c37c958121
Author: Xiang Fu <[email protected]>
AuthorDate: Wed Sep 30 20:55:15 2026 +0530

    Share sorted selection setup and trim duplicate coverage
    
    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.
---
 .../StreamingSelectionOrderByCombineOperator.java  |  32 +-
 .../query/StreamingSelectionOrderByOperator.java   | 121 +++-----
 ...reamingSelectionOrderByCombineOperatorTest.java |  74 +----
 .../StreamingSelectionOrderByOperatorTest.java     | 342 ++-------------------
 4 files changed, 104 insertions(+), 465 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 7ecf0a8cb9f..34d90e2a02d 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
@@ -217,27 +217,17 @@ public class StreamingSelectionOrderByCombineOperator 
extends BaseStreamingCombi
   /// ASC, descending by the column max value for DESC. Cursors without a 
min/max are placed first because they must
   /// always be processed (mirrors 
[MinMaxValueBasedSelectionOrderByCombineOperator]).
   private void sortCursorsByMinMax() {
-    if (_asc) {
-      Arrays.sort(_sortedCursors, (o1, o2) -> {
-        if (o1._minValue == null) {
-          return o2._minValue == null ? 0 : -1;
-        }
-        if (o2._minValue == null) {
-          return 1;
-        }
-        return o1._minValue.compareTo(o2._minValue);
-      });
-    } else {
-      Arrays.sort(_sortedCursors, (o1, o2) -> {
-        if (o1._maxValue == null) {
-          return o2._maxValue == null ? 0 : -1;
-        }
-        if (o2._maxValue == null) {
-          return 1;
-        }
-        return o2._maxValue.compareTo(o1._maxValue);
-      });
-    }
+    Arrays.sort(_sortedCursors, (o1, o2) -> {
+      Comparable bound1 = _asc ? o1._minValue : o1._maxValue;
+      Comparable bound2 = _asc ? o2._minValue : o2._maxValue;
+      if (bound1 == null) {
+        return bound2 == null ? 0 : -1;
+      }
+      if (bound2 == null) {
+        return 1;
+      }
+      return _asc ? bound1.compareTo(bound2) : bound2.compareTo(bound1);
+    });
   }
 
   @Override
diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/operator/query/StreamingSelectionOrderByOperator.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/operator/query/StreamingSelectionOrderByOperator.java
index 58f5f37ef7b..3010b9606f6 100644
--- 
a/pinot-core/src/main/java/org/apache/pinot/core/operator/query/StreamingSelectionOrderByOperator.java
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/operator/query/StreamingSelectionOrderByOperator.java
@@ -34,6 +34,7 @@ import javax.annotation.Nullable;
 import org.apache.pinot.common.request.context.ExpressionContext;
 import org.apache.pinot.common.request.context.OrderByExpressionContext;
 import org.apache.pinot.common.utils.DataSchema;
+import org.apache.pinot.common.utils.DataSchema.ColumnDataType;
 import org.apache.pinot.common.utils.HashUtil;
 import org.apache.pinot.core.common.BlockValSet;
 import org.apache.pinot.core.common.Operator;
@@ -196,7 +197,7 @@ public class StreamingSelectionOrderByOperator extends 
BaseOperator<SelectionRes
       _phase2NumColumns = 0;
       _phase2DataSourceMap = null;
       // Single-phase: all output expressions are order-by expressions, so 
their types are known up front.
-      _dataSchema = buildSinglePhaseDataSchema();
+      _dataSchema = buildDataSchema(_projectOperator::getResultColumnContext);
     }
     _numPhase1Columns = _phase1Expressions.size();
   }
@@ -216,7 +217,7 @@ public class StreamingSelectionOrderByOperator extends 
BaseOperator<SelectionRes
       // never have to reconstruct a schema the segment already knows. 
Two-phase leaves _dataSchema null until the
       // first fetch, which never happens here, so build it from column 
metadata alone.
       if (_twoPhase && _dataSchema == null) {
-        _dataSchema = 
buildTwoPhaseDataSchema(this::resolveResultColumnContext);
+        _dataSchema = buildDataSchema(this::resolveResultColumnContext);
       }
       assert _dataSchema != null;
       return new SelectionResultsBlock(_dataSchema, List.of(), _comparator, 
_queryContext);
@@ -240,34 +241,20 @@ public class StreamingSelectionOrderByOperator extends 
BaseOperator<SelectionRes
     if (remaining <= 0) {
       return null;
     }
-    ValueBlock valueBlock = nextNonEmptyBlock();
-    if (valueBlock == null) {
+    if (!loadNextBlock()) {
       return null;
     }
-    int numDocsFetched = valueBlock.getNumDocs();
-    BlockValSet[] blockValSets = new BlockValSet[_numPhase1Columns];
-    for (int i = 0; i < _numPhase1Columns; i++) {
-      blockValSets[i] = valueBlock.getBlockValueSet(_phase1Expressions.get(i));
-    }
-    RowBasedBlockValueFetcher blockValueFetcher = new 
RowBasedBlockValueFetcher(blockValSets);
-    int[] docIds = _twoPhase ? valueBlock.getDocIds() : null;
-    RoaringBitmap[] nullBitmaps = null;
-    if (_nullHandlingEnabled) {
-      nullBitmaps = new RoaringBitmap[_numPhase1Columns];
-      for (int i = 0; i < _numPhase1Columns; i++) {
-        nullBitmaps[i] = blockValSets[i].getNullBitmap();
-      }
-    }
-    _numDocsScanned += numDocsFetched;
-    _numEntriesScannedPostFilter += (long) numDocsFetched * 
_projectOperator.getNumColumnsProjected();
-    reportScanCost(numDocsFetched, (long) numDocsFetched * 
_projectOperator.getNumColumnsProjected());
 
     // Rows arrive sorted; we only need the first 'remaining' of them globally.
-    int numRows = Math.min(numDocsFetched, remaining);
+    int numRows = Math.min(_currentNumDocs, remaining);
     List<Object[]> rows = new ArrayList<>(numRows);
     for (int i = 0; i < numRows; i++) {
-      rows.add(materializeRow(blockValueFetcher, docIds, nullBitmaps, i));
+      rows.add(materializeRow(_currentFetcher, _currentDocIds, 
_currentNullBitmaps, i));
     }
+    _currentBlock = null;
+    _currentFetcher = null;
+    _currentDocIds = null;
+    _currentNullBitmaps = null;
     _numRowsEmitted += rows.size();
     return rows;
   }
@@ -319,37 +306,42 @@ public class StreamingSelectionOrderByOperator extends 
BaseOperator<SelectionRes
   /// is exhausted.
   @Nullable
   private Object[] nextRow() {
-    if (_currentBlock == null || _currentPos >= _currentNumDocs) {
-      if (_projectExhausted) {
-        return null;
-      }
-      _currentBlock = nextNonEmptyBlock();
-      if (_currentBlock == null) {
-        _projectExhausted = true;
-        return null;
-      }
-      BlockValSet[] blockValSets = new BlockValSet[_numPhase1Columns];
-      for (int i = 0; i < _numPhase1Columns; i++) {
-        blockValSets[i] = 
_currentBlock.getBlockValueSet(_phase1Expressions.get(i));
-      }
-      _currentFetcher = new RowBasedBlockValueFetcher(blockValSets);
-      _currentNumDocs = _currentBlock.getNumDocs();
-      _currentDocIds = _twoPhase ? _currentBlock.getDocIds() : null;
-      if (_nullHandlingEnabled) {
-        _currentNullBitmaps = new RoaringBitmap[_numPhase1Columns];
-        for (int i = 0; i < _numPhase1Columns; i++) {
-          _currentNullBitmaps[i] = blockValSets[i].getNullBitmap();
-        }
-      }
-      _currentPos = 0;
-      _numDocsScanned += _currentNumDocs;
-      _numEntriesScannedPostFilter += (long) _currentNumDocs * 
_projectOperator.getNumColumnsProjected();
-      reportScanCost(_currentNumDocs, (long) _currentNumDocs * 
_projectOperator.getNumColumnsProjected());
+    if ((_currentBlock == null || _currentPos >= _currentNumDocs) && 
!loadNextBlock()) {
+      return null;
     }
     int rowId = _currentPos++;
     return materializeRow(_currentFetcher, _currentDocIds, 
_currentNullBitmaps, rowId);
   }
 
+  private boolean loadNextBlock() {
+    if (_projectExhausted) {
+      return false;
+    }
+    _currentBlock = nextNonEmptyBlock();
+    if (_currentBlock == null) {
+      _projectExhausted = true;
+      return false;
+    }
+    BlockValSet[] blockValSets = new BlockValSet[_numPhase1Columns];
+    for (int i = 0; i < _numPhase1Columns; i++) {
+      blockValSets[i] = 
_currentBlock.getBlockValueSet(_phase1Expressions.get(i));
+    }
+    _currentFetcher = new RowBasedBlockValueFetcher(blockValSets);
+    _currentNumDocs = _currentBlock.getNumDocs();
+    _currentDocIds = _twoPhase ? _currentBlock.getDocIds() : null;
+    if (_nullHandlingEnabled) {
+      _currentNullBitmaps = new RoaringBitmap[_numPhase1Columns];
+      for (int i = 0; i < _numPhase1Columns; i++) {
+        _currentNullBitmaps[i] = blockValSets[i].getNullBitmap();
+      }
+    }
+    _currentPos = 0;
+    _numDocsScanned += _currentNumDocs;
+    _numEntriesScannedPostFilter += (long) _currentNumDocs * 
_projectOperator.getNumColumnsProjected();
+    reportScanCost(_currentNumDocs, (long) _currentNumDocs * 
_projectOperator.getNumColumnsProjected());
+    return true;
+  }
+
   /// Pulls the next project block carrying documents, skipping any that carry 
none. Returns `null` only when
   /// the project operator is exhausted; an empty block means "nothing in this 
batch", not "end of segment", and
   /// treating one as exhaustion would truncate the scan.
@@ -461,22 +453,11 @@ public class StreamingSelectionOrderByOperator extends 
BaseOperator<SelectionRes
       }
 
       if (_dataSchema == null) {
-        _dataSchema = 
buildTwoPhaseDataSchema(transformOperator::getResultColumnContext);
+        _dataSchema = 
buildDataSchema(transformOperator::getResultColumnContext);
       }
     }
   }
 
-  private DataSchema buildSinglePhaseDataSchema() {
-    String[] columnNames = new String[_numExpressions];
-    DataSchema.ColumnDataType[] columnDataTypes = new 
DataSchema.ColumnDataType[_numExpressions];
-    for (int i = 0; i < _numExpressions; i++) {
-      columnNames[i] = _expressions.get(i).toString();
-      columnDataTypes[i] = 
DataSchema.ColumnDataType.fromDataType(_orderByColumnContexts[i].getDataType(),
-          _orderByColumnContexts[i].isSingleValue());
-    }
-    return new DataSchema(columnNames, columnDataTypes);
-  }
-
   /// Resolves a non-order-by expression's result type without building the 
phase-2 pipeline, for the zero-match case
   /// where there is nothing to fetch. This reproduces 
[TransformOperator#getResultColumnContext] exactly -- that
   /// method resolves against its project operator's source column contexts, 
which for phase 2 are
@@ -494,21 +475,15 @@ public class StreamingSelectionOrderByOperator extends 
BaseOperator<SelectionRes
         TransformFunctionFactory.get(expression, _phase2ColumnContextMap, 
_queryContext));
   }
 
-  private DataSchema buildTwoPhaseDataSchema(Function<ExpressionContext, 
ColumnContext> resultColumnContexts) {
-    int numNonOrderByExpressions = _nonOrderByExpressions.size();
+  private DataSchema buildDataSchema(Function<ExpressionContext, 
ColumnContext> resultColumnContexts) {
     String[] columnNames = new String[_numExpressions];
-    DataSchema.ColumnDataType[] columnDataTypes = new 
DataSchema.ColumnDataType[_numExpressions];
+    ColumnDataType[] columnDataTypes = new ColumnDataType[_numExpressions];
     for (int i = 0; i < _numExpressions; i++) {
       columnNames[i] = _expressions.get(i).toString();
-    }
-    for (int i = 0; i < _numOrderByExpressions; i++) {
-      columnDataTypes[i] = 
DataSchema.ColumnDataType.fromDataType(_orderByColumnContexts[i].getDataType(),
-          _orderByColumnContexts[i].isSingleValue());
-    }
-    for (int i = 0; i < numNonOrderByExpressions; i++) {
-      ColumnContext columnContext = 
resultColumnContexts.apply(_nonOrderByExpressions.get(i));
-      columnDataTypes[_numOrderByExpressions + i] =
-          DataSchema.ColumnDataType.fromDataType(columnContext.getDataType(), 
columnContext.isSingleValue());
+      ColumnContext columnContext = i < _numOrderByExpressions ? 
_orderByColumnContexts[i]
+          : resultColumnContexts.apply(_expressions.get(i));
+      columnDataTypes[i] = 
ColumnDataType.fromDataType(columnContext.getDataType(),
+          columnContext.isSingleValue());
     }
     return new DataSchema(columnNames, columnDataTypes);
   }
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 2950d4062bf..0c050164188 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
@@ -521,35 +521,6 @@ public class StreamingSelectionOrderByCombineOperatorTest {
         "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
-  /// `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
@@ -1300,27 +1271,25 @@ public class 
StreamingSelectionOrderByCombineOperatorTest {
 
     int blockSize = 3;
     Result streamed = run(segments, query, true, nullHandling, true, 
blockSize);
-    assertStreamingParity(streamed, baseline, comparator, query, true, 
blockSize);
+    assertStreamingParity(streamed, baseline, comparator, query, blockSize);
   }
 
   private void assertStreamingParity(Result result, Result baseline, 
Comparator<Object[]> comparator,
-      @Language("sql") String query, boolean streaming, int blockSize) {
+      @Language("sql") String query, int blockSize) {
     assertEquals(result._combineOperator.getClass(), 
StreamingSelectionOrderByCombineOperator.class,
         "Expected the streaming combine operator for query: " + query);
     assertEquals(result._schema, baseline._schema, "Schema mismatch for query: 
" + query);
     assertSorted(result._rows, comparator);
     assertMultisetEquals(result._rows, baseline._rows);
-    if (streaming) {
-      int total = 0;
-      for (int size : result._blockSizes) {
-        assertTrue(size > 0 && size <= blockSize,
-            "Streamed block size out of range (0, " + blockSize + "] for query 
" + query + ": " + size);
-        total += size;
-      }
-      assertEquals(total, result._rows.size(), "Streamed block sizes must sum 
to the row count for query: " + query);
-      if (result._rows.size() > blockSize) {
-        assertTrue(result._numBlocks >= 2, "Expected multiple streamed blocks 
for query: " + query);
-      }
+    int total = 0;
+    for (int size : result._blockSizes) {
+      assertTrue(size > 0 && size <= blockSize,
+          "Streamed block size out of range (0, " + blockSize + "] for query " 
+ query + ": " + size);
+      total += size;
+    }
+    assertEquals(total, result._rows.size(), "Streamed block sizes must sum to 
the row count for query: " + query);
+    if (result._rows.size() > blockSize) {
+      assertTrue(result._numBlocks >= 2, "Expected multiple streamed blocks 
for query: " + query);
     }
   }
 
@@ -1485,27 +1454,6 @@ public class 
StreamingSelectionOrderByCombineOperatorTest {
         "The resolved mode must select the streaming combine");
   }
 
-  @Test
-  public void testAutoHonoursTheMinSortedRatioThresholdForDesc() {
-    // The DESC gate is a precondition, not a replacement for the ratio check: 
once reverse iteration is allowed, a
-    // DESC query must still clear the same threshold an ASC one does. 
_mixedSegments is 2 sorted of 4, so 0.5 passes
-    // and 0.75 does not, exactly as in 
testAutoHonoursTheMinSortedRatioThreshold.
-    String query = "SET allowReverseOrder=true; SELECT sortedCol, valCol FROM 
testTable ORDER BY sortedCol DESC, "
-        + "valCol DESC LIMIT 50";
-
-    QueryContext atThreshold =
-        modeContext("SET sortedSelectionMergeAutoMinSortedRatio=0.5; " + 
query, SortedSelectionMergeMode.AUTO);
-    makeStreamingInstancePlan(_mixedSegments, atThreshold);
-    assertEquals(atThreshold.getSortedSelectionMergeMode(), 
SortedSelectionMergeMode.ON,
-        "A DESC sorted ratio equal to the threshold must select the streaming 
merge");
-
-    QueryContext aboveThreshold =
-        modeContext("SET sortedSelectionMergeAutoMinSortedRatio=0.75; " + 
query, SortedSelectionMergeMode.AUTO);
-    makeStreamingInstancePlan(_mixedSegments, aboveThreshold);
-    assertEquals(aboveThreshold.getSortedSelectionMergeMode(), 
SortedSelectionMergeMode.OFF,
-        "A DESC sorted ratio below the threshold must keep the MinMax 
combine");
-  }
-
   @Test
   public void testAutoOnlyChecksTheLeadingOrderByDirection() {
     // Only the leading expression rides the segment's physical order, so a 
DESC tail is sorted in memory and needs no
diff --git 
a/pinot-core/src/test/java/org/apache/pinot/core/operator/query/StreamingSelectionOrderByOperatorTest.java
 
b/pinot-core/src/test/java/org/apache/pinot/core/operator/query/StreamingSelectionOrderByOperatorTest.java
index 62a83bc7816..d9b0719fb4d 100644
--- 
a/pinot-core/src/test/java/org/apache/pinot/core/operator/query/StreamingSelectionOrderByOperatorTest.java
+++ 
b/pinot-core/src/test/java/org/apache/pinot/core/operator/query/StreamingSelectionOrderByOperatorTest.java
@@ -20,21 +20,15 @@ package org.apache.pinot.core.operator.query;
 
 import java.io.File;
 import java.io.IOException;
-import java.util.ArrayDeque;
 import java.util.ArrayList;
-import java.util.Deque;
 import java.util.List;
-import java.util.Map;
 import java.util.concurrent.atomic.AtomicInteger;
 import org.apache.commons.io.FileUtils;
 import org.apache.pinot.common.request.context.ExpressionContext;
 import org.apache.pinot.common.request.context.OrderByExpressionContext;
 import org.apache.pinot.common.utils.DataSchema;
-import org.apache.pinot.core.common.BlockValSet;
 import org.apache.pinot.core.common.Operator;
 import org.apache.pinot.core.operator.BaseProjectOperator;
-import org.apache.pinot.core.operator.ColumnContext;
-import org.apache.pinot.core.operator.ExecutionStatistics;
 import org.apache.pinot.core.operator.blocks.ValueBlock;
 import org.apache.pinot.core.operator.blocks.results.SelectionResultsBlock;
 import org.apache.pinot.core.plan.ProjectPlanNode;
@@ -349,59 +343,20 @@ public class StreamingSelectionOrderByOperatorTest {
         "Expected no nulls in the output when null handling is disabled");
   }
 
-  /// No-tail mode, single phase: [StreamingSelectionOrderByOperator] must 
skip a zero-document project block
-  /// rather than read it as end-of-segment.
   @Test
-  public void testEmptyProjectBlockDoesNotTruncateSortedScan() {
-    assertEmptyProjectBlocksAreSkipped(_segment, "SELECT sortedCol FROM 
testTable ORDER BY sortedCol LIMIT 25", 1, 2,
-        1, 0);
-  }
-
-  /// Same, two phase: this is also the only path that asks the injected block 
for `getDocIds()`.
-  @Test
-  public void testEmptyProjectBlockDoesNotTruncateTwoPhaseSortedScan() {
-    assertEmptyProjectBlocksAreSkipped(_segment, "SELECT sortedCol, valCol 
FROM testTable ORDER BY sortedCol LIMIT 25",
-        1, 2, 1, 0);
-  }
-
-  /// Tail mode: the run scan reaches the project operator through 
`nextRow()`, which already tolerated empty
-  /// blocks. Pins that the two scan paths agree.
-  @Test
-  public void testEmptyProjectBlockDoesNotTruncateRunScan() {
-    assertEmptyProjectBlocksAreSkipped(_dupSegment,
-        "SELECT sortedCol, tailCol FROM testTable ORDER BY sortedCol, tailCol 
LIMIT 25", 1, 2, 1, 0);
-  }
-
-  /// The skip is a loop, not a single lookahead: several empty blocks in a 
row must all be skipped.
-  @Test
-  public void testConsecutiveEmptyProjectBlocksAreAllSkipped() {
-    assertEmptyProjectBlocksAreSkipped(_segment, "SELECT sortedCol FROM 
testTable ORDER BY sortedCol LIMIT 25", 1, 2,
-        3, 0);
-  }
-
-  /// Empty blocks arriving immediately before exhaustion exercise the loop's 
other exit: the skip must fall through
-  /// to the project operator's `null` rather than spin or emit a phantom row. 
The limit exceeds the fixture so
-  /// the scan actually reaches the end instead of stopping on the row budget.
-  @Test
-  public void testEmptyProjectBlocksBeforeExhaustionEndTheScan() {
-    assertEmptyProjectBlocksAreSkipped(_segment, "SELECT sortedCol FROM 
testTable ORDER BY sortedCol LIMIT 100", 1, 0,
-        0, 2);
-  }
-
-  /// The skip loop sits directly on top of the `limit + offset` budget 
bookkeeping, so cover a non-zero offset.
-  @Test
-  public void testEmptyProjectBlockIsSkippedWithOffset() {
-    assertEmptyProjectBlocksAreSkipped(_segment,
-        "SELECT sortedCol FROM testTable ORDER BY sortedCol LIMIT 15 OFFSET 
10", 1, 2, 1, 0);
-  }
-
-  /// `nextRow()` rebuilds the phase-1 null bitmaps in the same branch that 
pulls the next non-empty block, so
-  /// the skip must not desynchronise a bitmap from the block it was built 
against. Uses the fixture that carries a
-  /// null in the order-by tail column, with null handling on.
-  @Test
-  public void testEmptyProjectBlockIsSkippedWithNullHandling() {
-    assertEmptyProjectBlocksAreSkipped(_nullTailSegment,
-        "SELECT tailCol, sortedCol, valCol FROM testTable ORDER BY sortedCol, 
tailCol LIMIT 25", 1, true, 2, 1, 0);
+  public void testSkipsConsecutiveAndTrailingEmptyProjectBlocks() {
+    for (boolean tailToSort : new boolean[]{false, true}) {
+      IndexSegment segment = tailToSort ? _nullTailSegment : _segment;
+      String query = tailToSort
+          ? "SELECT tailCol, sortedCol, valCol FROM testTable ORDER BY 
sortedCol, tailCol LIMIT 100"
+          : "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol LIMIT 
100";
+      List<Object[]> expected = drain(buildOperator(segment, query, 1, true, 
false));
+      List<Object[]> actual = drain(buildOperator(segment, query, 1, true, 
true));
+      assertEquals(actual.size(), expected.size());
+      for (int i = 0; i < expected.size(); i++) {
+        assertEquals(actual.get(i), expected.get(i));
+      }
+    }
   }
 
   @Test
@@ -412,39 +367,8 @@ public class StreamingSelectionOrderByOperatorTest {
 
   @Test
   public void testZeroMatchSegmentEmitsOneSchemaBlockTwoPhase() {
-    assertZeroMatchEmitsOneSchemaBlock(_segment, "SELECT sortedCol, valCol 
FROM testTable WHERE sortedCol < 0 ORDER BY "
-        + "sortedCol", "SELECT sortedCol, valCol FROM testTable ORDER BY 
sortedCol", 1, false);
-  }
-
-  @Test
-  public void testZeroMatchSegmentEmitsOneSchemaBlockWithTailToSort() {
-    assertZeroMatchEmitsOneSchemaBlock(_dupSegment, "SELECT sortedCol, tailCol 
FROM testTable WHERE sortedCol < 0 "
-        + "ORDER BY sortedCol, tailCol", "SELECT sortedCol, tailCol FROM 
testTable ORDER BY sortedCol, tailCol", 1,
-        false);
-  }
-
-  @Test
-  public void testZeroMatchSegmentEmitsOneSchemaBlockWithNullHandling() {
-    assertZeroMatchEmitsOneSchemaBlock(_nullTailSegment, "SELECT tailCol, 
sortedCol, valCol FROM testTable WHERE "
-            + "sortedCol < 0 ORDER BY sortedCol, tailCol",
-        "SELECT tailCol, sortedCol, valCol FROM testTable ORDER BY sortedCol, 
tailCol", 1, true);
-  }
-
-  @Test
-  public void testMatchingSegmentEmitsNoTrailingEmptyBlock() {
-    assertNoEmptyBlockAmongEmittedBlocks(_segment, "SELECT sortedCol, valCol 
FROM testTable ORDER BY sortedCol", 1,
-        false);
-  }
-
-  @Test
-  public void testMatchingSegmentEmitsNoTrailingEmptyBlockWithTailToSort() {
-    assertNoEmptyBlockAmongEmittedBlocks(_dupSegment,
-        "SELECT sortedCol, tailCol, valCol FROM testTable ORDER BY sortedCol, 
tailCol", 1, false);
-  }
-
-  @Test
-  public void testMatchingSegmentEmitsNoTrailingEmptyBlockWithNullHandling() {
-    assertNoEmptyBlockAmongEmittedBlocks(_nullTailSegment,
+    assertZeroMatchEmitsOneSchemaBlock(_nullTailSegment,
+        "SELECT tailCol, sortedCol, valCol FROM testTable WHERE sortedCol < 0 
ORDER BY sortedCol, tailCol",
         "SELECT tailCol, sortedCol, valCol FROM testTable ORDER BY sortedCol, 
tailCol", 1, true);
   }
 
@@ -524,58 +448,13 @@ public class StreamingSelectionOrderByOperatorTest {
     return operator;
   }
 
-  /// Every block the tail-to-sort path emits carries a structurally mutable 
row list. Defensive rather than a fix for
-  /// a reachable failure: no consumer of this operator adds to the list 
today, but [SelectionResultsBlock] is a
-  /// shared type and `SelectionOperatorUtils.mergeWithoutOrdering()` adds to 
the row list of the block it merges
-  /// into, which a fixed-size `Arrays.asList` view would reject.
-  @Test
-  public void testRunPathEmitsGrowableRowLists() {
-    QueryContext queryContext = QueryContextConverterUtils.getQueryContext(
-        "SELECT sortedCol, tailCol FROM testTable ORDER BY sortedCol, 
tailCol");
-    queryContext.setSortedSelectionMergeMode(SortedSelectionMergeMode.ON);
-    Operator<SelectionResultsBlock> operator =
-        new SelectionPlanNode(new SegmentContext(_dupSegment), 
queryContext).run();
-    assertTrue(operator instanceof StreamingSelectionOrderByOperator,
-        "Expected the streaming operator, got: " + 
operator.getClass().getSimpleName());
-
-    int numBlocks = 0;
-    SelectionResultsBlock block;
-    while ((block = operator.nextBlock()) != null) {
-      numBlocks++;
-      List<Object[]> rows = block.getRows();
-      assertFalse(rows.isEmpty(), "Block " + numBlocks + " must carry rows, 
else the assertion below is vacuous");
-      int numRows = rows.size();
-      rows.add(new Object[]{0, 0});
-      assertEquals(rows.size(), numRows + 1, "Block " + numBlocks + " did not 
accept an appended row");
-    }
-    // One block per run, so more than one block proves the tail-to-sort path 
(and therefore drainAscending) ran; the
-    // no-tail path emits a single block for this fixture.
-    assertTrue(numBlocks > 1, "Fixture must exercise the tail-to-sort run 
path, got " + numBlocks + " block(s)");
-  }
-
-  /// The counterpart to the zero-match tests: a segment that does match rows 
must never emit an empty block, so the
-  /// guarantee reads "at least one block" and not "always an extra one".
-  private void assertNoEmptyBlockAmongEmittedBlocks(IndexSegment segment, 
@Language("sql") String query,
-      int numSortedExpressions, boolean nullHandling) {
-    StreamingSelectionOrderByOperator operator =
-        buildOperatorWithInjectableProject(segment, query, 
numSortedExpressions, nullHandling, 0, 0, 0);
-    int numBlocks = 0;
-    SelectionResultsBlock block;
-    while ((block = operator.nextBlock()) != null) {
-      numBlocks++;
-      assertFalse(block.getRows().isEmpty(),
-          "Block " + numBlocks + " is empty: the schema block must only be 
emitted when no rows were emitted at all");
-    }
-    assertTrue(numBlocks > 1, "Fixture must span more than one block, 
otherwise a trailing block cannot be observed");
-  }
-
   /// A segment matching no rows must still emit exactly one block before 
signalling exhaustion: empty, but carrying
   /// the same [DataSchema] the identical query produces when it does match 
rows. Consumers therefore never have
   /// to reconstruct a schema the segment already knows.
   private void assertZeroMatchEmitsOneSchemaBlock(IndexSegment segment, 
@Language("sql") String zeroMatchQuery,
       @Language("sql") String matchingQuery, int numSortedExpressions, boolean 
nullHandling) {
     SelectionResultsBlock matchingBlock =
-        buildOperatorWithInjectableProject(segment, matchingQuery, 
numSortedExpressions, nullHandling, 0, 0, 0)
+        buildOperator(segment, matchingQuery, numSortedExpressions, 
nullHandling, false)
             .nextBlock();
     assertNotNull(matchingBlock);
     assertFalse(matchingBlock.getRows().isEmpty(), "The control query must 
actually match rows");
@@ -583,7 +462,7 @@ public class StreamingSelectionOrderByOperatorTest {
     assertNotNull(expectedSchema);
 
     StreamingSelectionOrderByOperator operator =
-        buildOperatorWithInjectableProject(segment, zeroMatchQuery, 
numSortedExpressions, nullHandling, 0, 0, 0);
+        buildOperator(segment, zeroMatchQuery, numSortedExpressions, 
nullHandling, false);
     SelectionResultsBlock block = operator.nextBlock();
     assertNotNull(block, "A zero-match segment must still emit one block 
carrying the schema");
     assertTrue(block.getRows().isEmpty(), "A zero-match segment must not emit 
rows");
@@ -592,44 +471,8 @@ public class StreamingSelectionOrderByOperatorTest {
     assertNull(operator.nextBlock(), "Exactly one block may precede 
exhaustion");
   }
 
-  /// Drives the operator twice over `segment`: once against the real project 
operator, once against one that
-  /// splices zero-document blocks into its output, and asserts the emitted 
rows are identical.
-  ///
-  /// No doc-id-set operator produces an empty block today, so the case has to 
be injected. The project operator is
-  /// built with a deliberately small `maxDocsPerCall` so the scan spans 
several blocks.
-  ///
-  /// @param injectBeforeRealBlock 1-based index of the real block to splice 
empties in front of; 0 for none
-  /// @param numEmptyBlocks how many consecutive empty blocks to splice in at 
that point
-  /// @param numTrailingEmptyBlocks how many empty blocks to emit after the 
last real block, before exhaustion
-  private void assertEmptyProjectBlocksAreSkipped(IndexSegment segment, 
@Language("sql") String query,
-      int numSortedExpressions, int injectBeforeRealBlock, int numEmptyBlocks, 
int numTrailingEmptyBlocks) {
-    assertEmptyProjectBlocksAreSkipped(segment, query, numSortedExpressions, 
false, injectBeforeRealBlock,
-        numEmptyBlocks, numTrailingEmptyBlocks);
-  }
-
-  /// As above, with explicit control over null handling.
-  private void assertEmptyProjectBlocksAreSkipped(IndexSegment segment, 
@Language("sql") String query,
-      int numSortedExpressions, boolean nullHandling, int 
injectBeforeRealBlock, int numEmptyBlocks,
-      int numTrailingEmptyBlocks) {
-    List<Object[]> expected =
-        drain(buildOperatorWithInjectableProject(segment, query, 
numSortedExpressions, nullHandling, 0, 0, 0));
-    assertTrue(expected.size() > MAX_DOCS_PER_PROJECT_BLOCK,
-        "Fixture must span more than one project block, otherwise the 
injection is not exercised");
-    List<Object[]> actual = drain(
-        buildOperatorWithInjectableProject(segment, query, 
numSortedExpressions, nullHandling, injectBeforeRealBlock,
-            numEmptyBlocks, numTrailingEmptyBlocks));
-    assertEquals(actual.size(), expected.size(),
-        "Row count changed when empty project blocks were injected, so the 
scan was truncated");
-    for (int i = 0; i < expected.size(); i++) {
-      assertEquals(actual.get(i), expected.get(i), "Row " + i + " mismatch for 
query: " + query);
-    }
-  }
-
-  /// Builds a [StreamingSelectionOrderByOperator] directly (rather than 
through [SelectionPlanNode]) so
-  /// the project operator can be wrapped. Passing 0 for every injection 
parameter yields the undecorated operator.
-  private StreamingSelectionOrderByOperator 
buildOperatorWithInjectableProject(IndexSegment segment,
-      @Language("sql") String query, int numSortedExpressions, boolean 
nullHandling, int injectBeforeRealBlock,
-      int numEmptyBlocks, int numTrailingEmptyBlocks) {
+  private StreamingSelectionOrderByOperator buildOperator(IndexSegment segment,
+      @Language("sql") String query, int numSortedExpressions, boolean 
nullHandling, boolean injectEmptyBlocks) {
     QueryContext queryContext = 
QueryContextConverterUtils.getQueryContext(query);
     queryContext.setNullHandlingEnabled(nullHandling);
     queryContext.setSortedSelectionMergeMode(SortedSelectionMergeMode.ON);
@@ -647,11 +490,20 @@ public class StreamingSelectionOrderByOperatorTest {
     BaseProjectOperator<?> projectOperator =
         new ProjectPlanNode(new SegmentContext(segment), queryContext, 
projectExpressions,
             MAX_DOCS_PER_PROJECT_BLOCK).run();
-    BaseProjectOperator<?> effective =
-        injectBeforeRealBlock == 0 && numTrailingEmptyBlocks == 0
-            ? projectOperator
-            : new EmptyBlockInjectingProjectOperator(projectOperator, 
injectBeforeRealBlock, numEmptyBlocks,
-                numTrailingEmptyBlocks);
+    BaseProjectOperator<?> effective = projectOperator;
+    if (injectEmptyBlocks) {
+      effective = mock(BaseProjectOperator.class, 
delegatesTo(projectOperator));
+      ValueBlock emptyBlock = mock(ValueBlock.class);
+      AtomicInteger leading = new AtomicInteger();
+      AtomicInteger trailing = new AtomicInteger();
+      doAnswer(invocation -> {
+        if (leading.getAndIncrement() < 2) {
+          return emptyBlock;
+        }
+        ValueBlock block = projectOperator.nextBlock();
+        return block == null && trailing.getAndIncrement() < 2 ? emptyBlock : 
block;
+      }).when(effective).nextBlock();
+    }
     return new StreamingSelectionOrderByOperator(segment, queryContext, 
expressions, effective, numSortedExpressions);
   }
 
@@ -664,133 +516,6 @@ public class StreamingSelectionOrderByOperatorTest {
     return rows;
   }
 
-  /// Forwards everything to a real project operator, but splices 
zero-document blocks into its output:
-  /// `numEmptyBlocks` of them just before the `injectBeforeRealBlock`-th real 
block, and
-  /// `numTrailingEmptyBlocks` after the last real block but before exhaustion.
-  ///
-  /// Every real block is still delivered, merely deferred, so the decorator 
provably drops nothing of its own -
-  /// any row loss observed by a test is the operator under test truncating 
its scan.
-  private static class EmptyBlockInjectingProjectOperator extends 
BaseProjectOperator<ValueBlock> {
-    private final BaseProjectOperator<?> _delegate;
-    private final int _injectBeforeRealBlock;
-    private final int _numEmptyBlocks;
-    private final int _numTrailingEmptyBlocks;
-    private final Deque<ValueBlock> _pending = new ArrayDeque<>();
-    private int _realBlocksSeen;
-    private ValueBlock _lastRealBlock;
-    private boolean _injected;
-    private boolean _trailingEmitted;
-
-    EmptyBlockInjectingProjectOperator(BaseProjectOperator<?> delegate, int 
injectBeforeRealBlock, int numEmptyBlocks,
-        int numTrailingEmptyBlocks) {
-      _delegate = delegate;
-      _injectBeforeRealBlock = injectBeforeRealBlock;
-      _numEmptyBlocks = numEmptyBlocks;
-      _numTrailingEmptyBlocks = numTrailingEmptyBlocks;
-    }
-
-    @Override
-    protected ValueBlock getNextBlock() {
-      if (!_pending.isEmpty()) {
-        return _pending.poll();
-      }
-      ValueBlock block = _delegate.nextBlock();
-      if (block == null) {
-        // The empty blocks shadow the last real block so their value sets 
stay usable; there is nothing to shadow if
-        // the delegate never produced one, in which case the stream simply 
ends.
-        if (_trailingEmitted || _numTrailingEmptyBlocks == 0 || _lastRealBlock 
== null) {
-          return null;
-        }
-        _trailingEmitted = true;
-        for (int i = 0; i < _numTrailingEmptyBlocks; i++) {
-          _pending.add(new EmptyValueBlock(_lastRealBlock));
-        }
-        return _pending.poll();
-      }
-      _lastRealBlock = block;
-      _realBlocksSeen++;
-      if (!_injected && _realBlocksSeen == _injectBeforeRealBlock) {
-        _injected = true;
-        for (int i = 0; i < _numEmptyBlocks; i++) {
-          _pending.add(new EmptyValueBlock(block));
-        }
-        _pending.add(block);
-        return _pending.poll();
-      }
-      return block;
-    }
-
-    @Override
-    public Map<String, ColumnContext> getSourceColumnContextMap() {
-      return _delegate.getSourceColumnContextMap();
-    }
-
-    @Override
-    public ColumnContext getResultColumnContext(ExpressionContext expression) {
-      return _delegate.getResultColumnContext(expression);
-    }
-
-    @Override
-    public BaseProjectOperator<ValueBlock> withOrder(DocIdOrder newOrder) {
-      throw new UnsupportedOperationException("Test decorator does not support 
reordering");
-    }
-
-    @Override
-    public boolean isCompatibleWith(DocIdOrder order) {
-      return _delegate.isCompatibleWith(order);
-    }
-
-    @Override
-    public ExecutionStatistics getExecutionStatistics() {
-      return _delegate.getExecutionStatistics();
-    }
-
-    @Override
-    public List<? extends Operator> getChildOperators() {
-      return List.of(_delegate);
-    }
-
-    @Override
-    public String toExplainString() {
-      return "EMPTY_BLOCK_INJECTING_PROJECT";
-    }
-  }
-
-  /// A block reporting zero documents. Value sets are forwarded to the real 
block it shadows so that a consumer which
-  /// builds its fetchers before checking the document count does not fail for 
the wrong reason.
-  private static class EmptyValueBlock implements ValueBlock {
-    private final ValueBlock _delegate;
-
-    EmptyValueBlock(ValueBlock delegate) {
-      _delegate = delegate;
-    }
-
-    @Override
-    public int getNumDocs() {
-      return 0;
-    }
-
-    @Override
-    public int[] getDocIds() {
-      return new int[0];
-    }
-
-    @Override
-    public BlockValSet getBlockValueSet(ExpressionContext expression) {
-      return _delegate.getBlockValueSet(expression);
-    }
-
-    @Override
-    public BlockValSet getBlockValueSet(String column) {
-      return _delegate.getBlockValueSet(column);
-    }
-
-    @Override
-    public BlockValSet getBlockValueSet(String[] paths) {
-      return _delegate.getBlockValueSet(paths);
-    }
-  }
-
   /// Runs `query` twice over `segment` - once with the streaming hint on, 
once off - and asserts the
   /// concatenated streaming blocks equal the materialized operator's single 
block, cell by cell. Returns the (now
   /// exhausted) streaming operator so callers can make extra assertions on 
its execution statistics.
@@ -815,6 +540,7 @@ public class StreamingSelectionOrderByOperatorTest {
     SelectionResultsBlock block;
     while ((block = streamingOperator.nextBlock()) != null) {
       numBlocks++;
+      assertFalse(block.getRows().isEmpty(), "A matching segment must not emit 
empty blocks");
       if (streamingSchema == null) {
         streamingSchema = block.getDataSchema();
       }


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

Reply via email to