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 7efe592da816122a337106b02b4f010617357f45 Author: rohity <[email protected]> AuthorDate: Mon Sep 28 04:20:12 2026 +0000 Keep frontier pruning for segments flagged non-null. A query-level null-handling switch turned pruning off for every segment. The non-null metadata flag is already loaded with the segment, so flagged columns keep their min/max and unflagged ones always activate. AUTO uses the same predicate and no longer selects the streaming merge over segments the leaf will not stream. A consuming segment has no column metadata map, so it counts as unflagged rather than failing the query. --- .../StreamingSelectionOrderByCombineOperator.java | 21 +- .../apache/pinot/core/plan/SelectionPlanNode.java | 20 +- .../core/plan/maker/InstancePlanMakerImplV2.java | 30 +-- ...reamingSelectionOrderByCombineOperatorTest.java | 212 +++++++++++++++++++-- 4 files changed, 250 insertions(+), 33 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 5597a06bcfc..7ecf0a8cb9f 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 @@ -38,6 +38,7 @@ import org.apache.pinot.core.operator.blocks.results.MetadataResultsBlock; import org.apache.pinot.core.operator.blocks.results.SelectionResultsBlock; import org.apache.pinot.core.operator.query.StreamingSelectionOrderByOperator; import org.apache.pinot.core.operator.streaming.BaseStreamingCombineOperator; +import org.apache.pinot.core.plan.SelectionPlanNode; import org.apache.pinot.core.query.request.context.QueryContext; import org.apache.pinot.core.query.selection.SelectionOperatorUtils; import org.apache.pinot.core.query.utils.OrderByComparatorFactory; @@ -73,8 +74,11 @@ import org.slf4j.LoggerFactory; /// block read) when the merge frontier reaches its min/max, so once `limit + offset` rows are emitted the /// remaining segments are never acquired or read. See [#activateEligibleCursors(SegmentCursor)] for the /// correctness argument. -/// Pruning is disabled when null handling is enabled (an unsorted segment's first order-by column may then contain -/// nulls whose ordering position the raw min/max cannot capture), in which case every segment is activated. +/// With null handling on, a segment whose leading column is not flagged non-null by +/// [SelectionPlanNode#isColumnFlaggedNonNull] (it has nulls, the segment predates the flag, or it is a consuming +/// segment) is given no min/max and always activates. Flagged segments keep their min/max. With +/// null handling off, min/max are used as stored. The flag is read from segment metadata, not the null value vector: +/// that vector is a mapped buffer, and probing it here would acquire segments pruning is about to skip. /// Its effect shows in the stats: a never-activated segment scans no docs, so `numSegmentsMatched` excludes it while /// `numSegmentsProcessed` counts every segment. Their difference also counts activated segments that matched no rows. /// TODO: report the activated count separately; that needs a new `DataTable.MetadataKey`. @@ -165,7 +169,7 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi String firstOrderByColumn = firstOrderByExpressionContext.getType() == ExpressionContext.Type.IDENTIFIER ? firstOrderByExpressionContext.getIdentifier() : null; - _pruningEnabled = firstOrderByColumn != null && !queryContext.isNullHandlingEnabled(); + _pruningEnabled = firstOrderByColumn != null; // 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). @@ -185,7 +189,16 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi DataSourceMetadata metadata = operator.getIndexSegment().getDataSource(firstOrderByColumn, queryContext.getSchema()) .getDataSourceMetadata(); - _sortedCursors[i] = new SegmentCursor(operator, metadata.getMinValue(), metadata.getMaxValue()); + Comparable minValue = metadata.getMinValue(); + Comparable maxValue = metadata.getMaxValue(); + // isNonNull is loaded with the segment, the same class of access as min/max. False means nulls or unknown, so + // drop the bounds and let activateEligibleCursors force-activate this cursor. Do not read the null vector here. + if (queryContext.isNullHandlingEnabled() + && !SelectionPlanNode.isColumnFlaggedNonNull(operator.getIndexSegment(), firstOrderByColumn)) { + minValue = null; + maxValue = null; + } + _sortedCursors[i] = new SegmentCursor(operator, minValue, maxValue); if (metadata.isSorted()) { numSortedSegments++; } diff --git a/pinot-core/src/main/java/org/apache/pinot/core/plan/SelectionPlanNode.java b/pinot-core/src/main/java/org/apache/pinot/core/plan/SelectionPlanNode.java index f2a39b14777..c022bd76173 100644 --- a/pinot-core/src/main/java/org/apache/pinot/core/plan/SelectionPlanNode.java +++ b/pinot-core/src/main/java/org/apache/pinot/core/plan/SelectionPlanNode.java @@ -20,6 +20,7 @@ package org.apache.pinot.core.plan; import java.util.ArrayList; import java.util.List; +import java.util.Map; import org.apache.pinot.common.request.context.ExpressionContext; import org.apache.pinot.common.request.context.OrderByExpressionContext; import org.apache.pinot.common.utils.config.QueryOptionsUtils; @@ -35,6 +36,7 @@ import org.apache.pinot.core.operator.query.SelectionPartiallyOrderedByLinearOpe import org.apache.pinot.core.operator.query.StreamingSelectionOrderByOperator; import org.apache.pinot.core.query.request.context.QueryContext; import org.apache.pinot.core.query.selection.SelectionOperatorUtils; +import org.apache.pinot.segment.spi.ColumnMetadata; import org.apache.pinot.segment.spi.IndexSegment; import org.apache.pinot.segment.spi.SegmentContext; import org.apache.pinot.segment.spi.datasource.DataSource; @@ -231,7 +233,9 @@ public class SelectionPlanNode implements PlanNode { /// This is [#isColumnPhysicallySorted] plus the null caveat: once null handling is on, the physical order is not /// the order the query asks for, so a column carrying nulls is not usable as sorted. Reading the null bitmap /// touches the segment's mapped buffer, so this must only be called once the segment has been acquired -- which is - /// why [AcquireReleaseColumnsSegmentPlanNode] defers the whole plan build until after `acquire()`. + /// why [AcquireReleaseColumnsSegmentPlanNode] defers the whole plan build until after `acquire()`. Callers that + /// run before acquire use [#isColumnFlaggedNonNull] instead. That flag is one-sided (`true` means no nulls; `false` + /// means nulls or unknown) and is stricter than the bitmap check below. public static boolean isColumnSorted(IndexSegment segment, QueryContext queryContext, String column) { DataSource dataSource = segment.getDataSource(column, queryContext.getSchema()); // If there are null values, we cannot trust DataSourceMetadata.isSorted @@ -253,4 +257,18 @@ public class SelectionPlanNode implements PlanNode { public static boolean isColumnPhysicallySorted(IndexSegment segment, QueryContext queryContext, String column) { return segment.getDataSource(column, queryContext.getSchema()).getDataSourceMetadata().isSorted(); } + + /// Returns whether segment metadata flags `column` as holding no nulls ([ColumnMetadata#isNonNull]). + /// + /// Reads segment metadata only, so it needs no segment acquire. `false` means nulls or unknown. A consuming segment + /// has no column metadata map (it is `null`, so + /// [org.apache.pinot.segment.spi.SegmentMetadata#getColumnMetadataFor] would throw) and is never flagged. + public static boolean isColumnFlaggedNonNull(IndexSegment segment, String column) { + Map<String, ColumnMetadata> columnMetadataMap = segment.getSegmentMetadata().getColumnMetadataMap(); + if (columnMetadataMap == null) { + return false; + } + ColumnMetadata columnMetadata = columnMetadataMap.get(column); + return columnMetadata != null && columnMetadata.isNonNull(); + } } diff --git a/pinot-core/src/main/java/org/apache/pinot/core/plan/maker/InstancePlanMakerImplV2.java b/pinot-core/src/main/java/org/apache/pinot/core/plan/maker/InstancePlanMakerImplV2.java index f00ffb4aa18..6b7aa155c98 100644 --- a/pinot-core/src/main/java/org/apache/pinot/core/plan/maker/InstancePlanMakerImplV2.java +++ b/pinot-core/src/main/java/org/apache/pinot/core/plan/maker/InstancePlanMakerImplV2.java @@ -505,28 +505,30 @@ public class InstancePlanMakerImplV2 implements PlanMaker { // // TODO: allowReverseOrder=true still does not guarantee streaming leaves. getSortedByProject() swallows the // failure when a docId set cannot be reversed and returns the unreversed operator, so the leaf gate then - // declines and AUTO has again resolved ON over materialized children. Same residual shape as the null - // handling case below; a shared "can this segment actually stream in the query's direction" predicate - // would close both. See https://github.com/apache/pinot/pull/19120#discussion_r3871713989 + // declines and AUTO has resolved ON over materialized children. + // See https://github.com/apache/pinot/pull/19120#discussion_r3871713989 LOGGER.debug("Not using {} for a DESC leading ORDER BY expression without {}: {}", QueryOptionKey.SORTED_SELECTION_MERGE_MODE, QueryOptionKey.ALLOW_REVERSE_ORDER, firstOrderByExpression); return false; } String column = firstOrderByExpression.getIdentifier(); + boolean nullHandlingEnabled = queryContext.isNullHandlingEnabled(); int numSorted = 0; for (SegmentContext segmentContext : segmentContexts) { - // TODO: This deliberately uses the physical-sortedness predicate rather than the full - // SelectionPlanNode.isColumnSorted(), which additionally rejects a column carrying nulls when null - // handling is on. That null check reads the segment's mapped buffer via the null value vector, and this - // runs before any segment is acquired -- AcquireReleaseColumnsSegmentPlanNode defers the whole plan build - // precisely so that no planner touches a buffer pre-acquire. Consequence: with null handling on and a - // null-bearing leading column, AUTO can resolve to ON while the leaf gate then refuses to stream, giving - // a k-way merge over fully materialized children -- correct, but slower than the MinMax combine. Handle - // properly (an acquire-free nullability signal, or resolving AUTO per segment at acquire time) as a - // follow-up. See https://github.com/apache/pinot/pull/19120#discussion_r3871713975 - if (SelectionPlanNode.isColumnPhysicallySorted(segmentContext.getIndexSegment(), queryContext, column)) { - numSorted++; + // Physical sortedness is metadata. SelectionPlanNode.isColumnSorted() also rejects a null-bearing column, but + // that reads the null value vector, a mapped buffer, and this runs before any segment is acquired. Use + // SelectionPlanNode.isColumnFlaggedNonNull() instead. It is one-sided (false means nulls or unknown) and + // stricter than the leaf, which can see an empty bitmap after acquire: an old segment with no flag resolves OFF + // rather than streaming over a materialized child. + // See https://github.com/apache/pinot/pull/19120#discussion_r3871713975 + IndexSegment segment = segmentContext.getIndexSegment(); + if (!SelectionPlanNode.isColumnPhysicallySorted(segment, queryContext, column)) { + continue; } + if (nullHandlingEnabled && !SelectionPlanNode.isColumnFlaggedNonNull(segment, column)) { + continue; + } + numSorted++; } double sortedRatio = (double) numSorted / numSegments; double minSortedRatio = queryContext.getSortedSelectionMergeAutoMinSortedRatio(); 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 606705708cc..2532d4aacc7 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 @@ -24,6 +24,7 @@ import java.util.ArrayList; import java.util.Comparator; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.stream.Collectors; @@ -43,6 +44,7 @@ import org.apache.pinot.core.operator.query.StreamingSelectionOrderByOperator; import org.apache.pinot.core.plan.CombinePlanNode; import org.apache.pinot.core.plan.ExplainInfo; import org.apache.pinot.core.plan.PlanNode; +import org.apache.pinot.core.plan.SelectionPlanNode; import org.apache.pinot.core.plan.maker.InstancePlanMakerImplV2; import org.apache.pinot.core.plan.maker.PlanMaker; import org.apache.pinot.core.query.executor.ResultsBlockStreamer; @@ -51,8 +53,11 @@ import org.apache.pinot.core.query.request.context.utils.QueryContextConverterUt import org.apache.pinot.core.query.utils.OrderByComparatorFactory; import org.apache.pinot.core.util.QueryMultiThreadingUtils; import org.apache.pinot.segment.local.indexsegment.immutable.ImmutableSegmentLoader; +import org.apache.pinot.segment.local.indexsegment.mutable.MutableSegmentImpl; +import org.apache.pinot.segment.local.indexsegment.mutable.MutableSegmentImplTestUtils; import org.apache.pinot.segment.local.segment.creator.impl.SegmentIndexCreationDriverImpl; import org.apache.pinot.segment.local.segment.readers.GenericRowRecordReader; +import org.apache.pinot.segment.spi.ColumnMetadata; import org.apache.pinot.segment.spi.IndexSegment; import org.apache.pinot.segment.spi.SegmentContext; import org.apache.pinot.segment.spi.creator.SegmentGeneratorConfig; @@ -161,6 +166,14 @@ public class StreamingSelectionOrderByCombineOperatorTest { /// Two segments whose rows tie on every order-by expression; see [#buildTiedSortedRecords]. private List<IndexSegment> _tiedSegments; + /// Three disjoint ranges built with null handling and no nulls in `sortedCol`, so the column is flagged non-null. + /// Ranges are `[0, 100)`, `[1000, 1100)`, `[2000, 2100)`. + private List<IndexSegment> _nonNullDisjointSegments; + /// The middle of those ranges, with nulls in `sortedCol`. Stored nulls are `Integer.MIN_VALUE`, so the metadata min + /// is not an ASC bound; the metadata max is the highest non-null, which a DESC top-5 over the high range would prune + /// if the flag were ignored. + private IndexSegment _nullBearingMidSegment; + @BeforeClass public void setUp() throws Exception { @@ -210,6 +223,14 @@ public class StreamingSelectionOrderByCombineOperatorTest { _interleavedSegments = new ArrayList<>(2); _interleavedSegments.add(buildSegment(SORTED_TABLE_CONFIG, "sparse_0", buildSparseSortedRecords(), false)); _interleavedSegments.add(buildSegment(SORTED_TABLE_CONFIG, "denseGap_0", buildGapFillingSortedRecords(), false)); + + _nonNullDisjointSegments = new ArrayList<>(3); + for (int i = 0; i < 3; i++) { + _nonNullDisjointSegments.add( + buildSegment(SORTED_TABLE_CONFIG, "nonNullDisjoint_" + i, buildDisjointSortedRecords(i), true)); + } + _nullBearingMidSegment = buildSegment(SORTED_TABLE_CONFIG, "nullBearingMid_0", + buildNullBearingDisjointRecords(1), true); } /// Leading order-by values `0, 1000, 2000, ...`: gaps wide enough for another segment's entire range to sit between @@ -266,6 +287,17 @@ public class StreamingSelectionOrderByCombineOperatorTest { return records; } + /// Like [#buildDisjointSortedRecords] but the leading order-by column itself is null in the first few rows. + private static List<GenericRow> buildNullBearingDisjointRecords(int index) { + List<GenericRow> records = buildDisjointSortedRecords(index); + for (int i = 0; i < 3; i++) { + GenericRow record = records.get(i); + record.putValue(SORTED_COL, null); + record.addNullValueField(SORTED_COL); + } + return records; + } + /// Like [#buildOverlappingSortedRecords] but the leading order-by column itself is null in the first few rows. private static List<GenericRow> buildNullBearingSortedRecords(int index) { List<GenericRow> records = buildOverlappingSortedRecords(index); @@ -365,6 +397,32 @@ public class StreamingSelectionOrderByCombineOperatorTest { return ImmutableSegmentLoader.load(new File(TEMP_DIR, segmentName), ReadMode.mmap); } + /// [ColumnMetadata#isNonNull] is the acquire-free signal pruning and AUTO use. These fixtures have to carry it the + /// way [org.apache.pinot.segment.local.segment.creator.impl.BaseSegmentCreator] writes it, or the tests below are + /// vacuous. + @Test + public void testNonNullFlagMatchesHowSegmentsWereBuilt() { + for (IndexSegment segment : _sortedSegments) { + assertTrue(isColumnNonNull(segment, SORTED_COL), segment.getSegmentName()); + assertFalse(isColumnNonNull(segment, NULLABLE_COL), segment.getSegmentName()); + } + for (IndexSegment segment : _disjointSegments) { + assertFalse(isColumnNonNull(segment, SORTED_COL), segment.getSegmentName()); + } + for (IndexSegment segment : _nullBearingSortedSegments) { + assertFalse(isColumnNonNull(segment, SORTED_COL), segment.getSegmentName()); + } + for (IndexSegment segment : _nonNullDisjointSegments) { + assertTrue(isColumnNonNull(segment, SORTED_COL), segment.getSegmentName()); + } + assertFalse(isColumnNonNull(_nullBearingMidSegment, SORTED_COL)); + } + + private static boolean isColumnNonNull(IndexSegment segment, String column) { + ColumnMetadata metadata = segment.getSegmentMetadata().getColumnMetadataFor(column); + return metadata != null && metadata.isNonNull(); + } + @Test public void testAscendingParity() { assertParity(_sortedSegments, "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol LIMIT 50", false); @@ -587,7 +645,7 @@ public class StreamingSelectionOrderByCombineOperatorTest { @Test public void testNullHandlingEnabledParity() { - // Null handling on disables min/max pruning (the combine activates every segment); nullableCol carries real nulls. + // nullableCol carries real nulls. sortedCol is flagged non-null, so frontier pruning stays on for this order-by. assertParity(_sortedSegments, "SELECT nullableCol, sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol LIMIT 50", true); } @@ -635,6 +693,126 @@ public class StreamingSelectionOrderByCombineOperatorTest { assertEquals(result._numSegmentsMatched, 1, "numSegmentsMatched must exclude never-activated segments"); } + /// Null handling on, but every segment's leading column is flagged non-null, so min/max pruning stays exact. + @Test + public void testNullHandlingPrunesFlaggedNonNullSegments() { + @Language("sql") String query = + "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol LIMIT 5"; + assertParity(_nonNullDisjointSegments, query, true); + Result result = run(_nonNullDisjointSegments, query, true, true, true, 0); + assertEquals(result._rows.size(), 5); + assertSegmentsScanned(result, 1, _nonNullDisjointSegments.size()); + } + + /// These segments were built with null handling off, so `sortedCol` is unflagged even though it holds no nulls. + /// With null handling on the query, every segment must activate: an unknown flag is not proof of no nulls. + @Test + public void testUnflaggedSegmentsAlwaysActivateUnderNullHandling() { + @Language("sql") String query = + "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol LIMIT 5"; + Result result = run(_disjointSegments, query, true, true, true, 0); + assertEquals(result._rows.size(), 5); + assertSegmentsScanned(result, NUM_SEGMENTS, NUM_SEGMENTS); + } + + /// Flagged low and high ranges around one null-bearing middle segment. ASC top-5 lives entirely in the low range; + /// the high range prunes, and the null-bearing segment still activates. On DESC a null frontier head stops further + /// pruning, so the skipped-segment claim is made by acquiring the null-bearing segment rather than by a scan count. + @Test + public void testNullBearingSegmentForceActivatesAlongsideFlaggedOnes() { + List<IndexSegment> segments = nullPruningMix(); + @Language("sql") String query = + "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol LIMIT 5"; + assertParity(segments, query, true); + Result result = run(segments, query, true, true, true, 0); + assertEquals(result._rows.size(), 5); + for (int i = 0; i < 5; i++) { + assertEquals((int) result._rows.get(i)[0], i); + } + assertSegmentsScanned(result, 2, segments.size()); + // The ASC count above does not prove this segment was force-activated: its stored null is Integer.MIN_VALUE, so + // the metadata min would activate it anyway. On DESC the metadata max is the highest non-null and sits below the + // high range, so keeping that max would never acquire it. Its nulls sort first, which is what parity checks. + @Language("sql") String desc = + "SET allowReverseOrder=true; SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol DESC, valCol DESC " + + "LIMIT 5"; + assertParity(segments, desc, true); + QueryContext descContext = hintedContext(desc); + descContext.setNullHandlingEnabled(true); + List<InstrumentedSegmentOperator> children = instrument(segments, descContext); + drain(combineOver(children, descContext)); + assertTrue(children.get(1)._numAcquires > 0, + "The null-bearing segment must be acquired; its metadata max would prune it"); + } + + /// DESC with `allowReverseOrder`. The low segment was built with null handling off, so it is unflagged and its max + /// sits below the high range: trusting that max would prune it, and with it the reverse bitmap built on the first + /// block. The middle flagged segment stays past the frontier and is never scanned. A null-bearing segment cannot + /// stand in here, because a null frontier head refuses every later prune. + @Test + public void testDescReverseScanOnlyForceActivatesUnflaggedSegments() { + List<IndexSegment> segments = List.of(_disjointSegments.get(0), _nonNullDisjointSegments.get(1), + _nonNullDisjointSegments.get(2)); + @Language("sql") String query = + "SET allowReverseOrder=true; SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol DESC, valCol DESC " + + "LIMIT 5"; + assertParity(segments, query, true); + Result result = run(segments, query, true, true, true, 0); + assertEquals(result._rows.size(), 5); + int top = 2 * 1000 + NUM_RECORDS_PER_SEGMENT - 1; + for (int i = 0; i < 5; i++) { + assertEquals((int) result._rows.get(i)[0], top - i); + } + assertSegmentsScanned(result, 2, segments.size()); + } + + /// A consuming segment has no column metadata map, so reading its non-null flag through + /// `SegmentMetadata.getColumnMetadataFor` throws. It tracks min/max as it ingests, though, so if it were treated as + /// flagged its bounds would prune it here just like the high segment. It has to count as unflagged and activate. + @Test + public void testConsumingSegmentAlwaysActivatesUnderNullHandling() + throws Exception { + MutableSegmentImpl consumingSegment = + MutableSegmentImplTestUtils.createMutableSegmentImpl(SCHEMA, Set.of(), Set.of(), Set.of(), false, true); + try { + for (GenericRow record : buildDisjointSortedRecords(1)) { + consumingSegment.index(record, null); + } + assertNull(consumingSegment.getSegmentMetadata().getColumnMetadataMap(), + "The fixture must reproduce the consuming segment's missing column metadata, or this test is vacuous"); + assertFalse(SelectionPlanNode.isColumnFlaggedNonNull(consumingSegment, SORTED_COL)); + + List<IndexSegment> segments = + List.of(_nonNullDisjointSegments.get(0), consumingSegment, _nonNullDisjointSegments.get(2)); + @Language("sql") String query = "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol LIMIT 5"; + assertParity(segments, query, true); + Result result = run(segments, query, true, true, true, 0); + assertEquals(result._rows.size(), 5); + for (int i = 0; i < 5; i++) { + assertEquals((int) result._rows.get(i)[0], i); + } + assertSegmentsScanned(result, 2, segments.size()); + } finally { + consumingSegment.destroy(); + } + } + + /// Low flagged, null-bearing middle, high flagged. The middle segment is the only one that must force-activate. + private List<IndexSegment> nullPruningMix() { + return List.of(_nonNullDisjointSegments.get(0), _nullBearingMidSegment, _nonNullDisjointSegments.get(2)); + } + + /// `numSegmentsMatched` is the count of children that scanned a doc, which is the cursors the merge activated. A + /// pruned child scans nothing, which is also what keeps it from building a DESC reverse bitmap. Doc count is only + /// an upper bound: an activated streaming child may stop after the first rows of its block. + private static void assertSegmentsScanned(Result result, int numScanned, int numSegments) { + assertEquals(result._numSegmentsProcessed, numSegments, "numSegmentsProcessed counts every segment"); + assertEquals(result._numSegmentsMatched, numScanned, + "numSegmentsMatched counts segments that scanned a doc; scanned docs: " + result._numDocsScanned); + assertTrue(result._numDocsScanned > 0 && result._numDocsScanned <= (long) numScanned * NUM_RECORDS_PER_SEGMENT, + "Scanned " + result._numDocsScanned + " docs across " + numScanned + " segments"); + } + /// A cursor activated *late* can hold a smaller row than the retained leader. Two things must be right for it to /// land in the correct place: the leader is re-compared after `activateEligibleCursors()` rather than before, and /// the pruning frontier is the row about to be emitted rather than the last one emitted. Either mistake emits the @@ -1077,9 +1255,10 @@ public class StreamingSelectionOrderByCombineOperatorTest { public void testExplainAttributesReflectTheQueryShape() { assertTrue(explainAttributes(_sortedSegments, "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol LIMIT 5", false).get("deferTiedCursors").getBool(), "A single order-by expression defers ties"); - assertFalse(explainAttributes(_sortedSegments, + assertTrue(explainAttributes(_sortedSegments, "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol LIMIT 5", true) - .get("frontierPruning").getBool(), "Null handling disables frontier pruning"); + .get("frontierPruning").getBool(), + "Null handling no longer disables frontier pruning when the leading order-by is a column"); assertEquals(explainAttributes(_mixedSegments, "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol LIMIT 5", false) .get("numSortedSegments").getLong(), 2); @@ -1202,14 +1381,11 @@ public class StreamingSelectionOrderByCombineOperatorTest { } @Test - public void testAutoIgnoresNullsInTheLeadingColumn() { - // Pins the known gap documented by the TODO in InstancePlanMakerImplV2#isSortedEnoughForStreamingMerge, so it is - // enforced by CI rather than only described. AUTO resolves from physical sortedness alone, because the null check - // that SelectionPlanNode#isColumnSorted adds reads the segment's mapped buffer and AUTO runs before any acquire. - // So a null-bearing sorted column resolves to ON with null handling either off or on -- but with it on, the leaf - // gate then refuses to stream, and the merge runs over fully materialized children. Correct, just not fast. - // - // Tighten this to OFF for the null-handling case when the follow-up lands. + public void testAutoRequiresNonNullLeadingColumnWhenNullHandlingIsOn() { + // AUTO runs before any segment is acquired, so it must not read the null value vector. It uses + // ColumnMetadata.isNonNull() instead, which is one-sided: false means nulls or unknown. That is stricter than + // SelectionPlanNode.isColumnSorted(), which still reads the bitmap after acquire. An old segment with an empty + // bitmap and no flag therefore resolves OFF rather than streaming over a materialized child. // See https://github.com/apache/pinot/pull/19120#discussion_r3871713975 String query = "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol LIMIT 50"; @@ -1221,8 +1397,8 @@ public class StreamingSelectionOrderByCombineOperatorTest { QueryContext nullHandlingOn = modeContext(query, SortedSelectionMergeMode.AUTO); nullHandlingOn.setNullHandlingEnabled(true); makeStreamingInstancePlan(_nullBearingSortedSegments, nullHandlingOn); - assertEquals(nullHandlingOn.getSortedSelectionMergeMode(), SortedSelectionMergeMode.ON, - "AUTO currently reads physical sortedness only, so null handling does not change the decision"); + assertEquals(nullHandlingOn.getSortedSelectionMergeMode(), SortedSelectionMergeMode.OFF, + "With null handling on an unflagged leading column must not count as sorted"); // The leaf, which runs post-acquire and can afford the null check, still refuses to stream these segments. List<SegmentContext> segmentContexts = new ArrayList<>(1); segmentContexts.add(new SegmentContext(_nullBearingSortedSegments.get(0))); @@ -1230,6 +1406,12 @@ public class StreamingSelectionOrderByCombineOperatorTest { assertFalse(leafOperator instanceof StreamingSelectionOrderByOperator, "A null-bearing leading column must not produce a streaming leaf under null handling, got: " + leafOperator.getClass().getSimpleName()); + + QueryContext flagged = modeContext(query, SortedSelectionMergeMode.AUTO); + flagged.setNullHandlingEnabled(true); + makeStreamingInstancePlan(_sortedSegments, flagged); + assertEquals(flagged.getSortedSelectionMergeMode(), SortedSelectionMergeMode.ON, + "A leading column flagged non-null stays sorted under null handling"); } @Test @@ -1489,11 +1671,13 @@ public class StreamingSelectionOrderByCombineOperatorTest { throws IOException { EXECUTOR.shutdownNow(); for (List<IndexSegment> segments : List.of(_sortedSegments, _disjointSegments, _mixedSegments, _lowCardSegments, - _interleavedSegments, _tiedSegments)) { + _interleavedSegments, _tiedSegments, _nullBearingSortedSegments, _unsortedSegments, + _nonNullDisjointSegments)) { for (IndexSegment segment : segments) { segment.destroy(); } } + _nullBearingMidSegment.destroy(); FileUtils.deleteDirectory(TEMP_DIR); } --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
