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]
