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]

Reply via email to