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 849901a015daf0a479f8a281ca7b639e53e22078
Author: rohity <[email protected]>
AuthorDate: Wed Aug 5 18:24:30 2026 +0000

    Add streaming selection ORDER BY combine for physically sorted segments
    
    An unbounded leaf-stage ORDER BY (as injected for sorted merge join inputs)
    routes to MinMaxValueBasedSelectionOrderByCombineOperator, which merges 
every
    segment's rows into a single block before returning anything. At large data
    volumes this exceeds the leaf stage's CPU budget and ThreadAccountant raises
    EarlyTerminationException inside SelectionOperatorUtils.mergeWithOrdering(),
    surfacing to the broker as a spurious "Cancelled by sender".
    
    This adds a streaming alternative, opt-in via the 
`streamingSelectionOrderBy`
    query option:
    
    - StreamingSelectionOrderByOperator emits sorted blocks incrementally for a
      segment that is physically sorted on the leading ORDER BY column, reading
      the sorted forward index in order instead of building a priority queue.
    - StreamingSelectionOrderByCombineOperator performs a k-way heap merge 
across
      segment operators and emits bounded blocks 
(`streamingSelectionOrderByBlockSize`,
      default 10000) rather than one materialized result.
    - SelectionPlanNode and CombinePlanNode select these operators when the 
option
      is set and the sortedness precondition holds; otherwise behaviour is 
unchanged.
    
    Part of #18667.
---
 .../common/utils/config/QueryOptionsUtils.java     |  10 +
 .../common/utils/config/QueryOptionsUtilsTest.java |   4 +-
 .../StreamingSelectionOrderByCombineOperator.java  | 543 ++++++++++++++++++++
 .../query/StreamingSelectionOrderByOperator.java   | 514 +++++++++++++++++++
 .../apache/pinot/core/plan/CombinePlanNode.java    |  16 +
 .../apache/pinot/core/plan/SelectionPlanNode.java  |  28 +-
 .../core/plan/maker/InstancePlanMakerImplV2.java   |  11 +
 .../core/query/request/context/QueryContext.java   |  22 +
 .../operator/combine/CombineSlowOperatorsTest.java |  94 +++-
 ...reamingSelectionOrderByCombineOperatorTest.java | 553 +++++++++++++++++++++
 .../StreamingSelectionOrderByOperatorTest.java     | 406 +++++++++++++++
 .../apache/pinot/spi/utils/CommonConstants.java    |  10 +
 12 files changed, 2198 insertions(+), 13 deletions(-)

diff --git 
a/pinot-common/src/main/java/org/apache/pinot/common/utils/config/QueryOptionsUtils.java
 
b/pinot-common/src/main/java/org/apache/pinot/common/utils/config/QueryOptionsUtils.java
index b028106ed68..d8393720613 100644
--- 
a/pinot-common/src/main/java/org/apache/pinot/common/utils/config/QueryOptionsUtils.java
+++ 
b/pinot-common/src/main/java/org/apache/pinot/common/utils/config/QueryOptionsUtils.java
@@ -602,6 +602,16 @@ public class QueryOptionsUtils {
     return 
checkedParseIntNonNegative(QueryOptionKey.STREAMING_DISTINCT_FLUSH_THRESHOLD, 
value);
   }
 
+  public static boolean isSortedSelectionMergeEnabled(Map<String, String> 
queryOptions) {
+    return 
Boolean.parseBoolean(queryOptions.get(QueryOptionKey.SORTED_SELECTION_MERGE_ENABLED));
+  }
+
+  @Nullable
+  public static Integer getSortedSelectionMergeBlockSize(Map<String, String> 
queryOptions) {
+    String value = 
queryOptions.get(QueryOptionKey.SORTED_SELECTION_MERGE_BLOCK_SIZE);
+    return 
checkedParseIntPositive(QueryOptionKey.SORTED_SELECTION_MERGE_BLOCK_SIZE, 
value);
+  }
+
   public static boolean isNullHandlingEnabled(Map<String, String> 
queryOptions) {
     return 
Boolean.parseBoolean(queryOptions.get(QueryOptionKey.ENABLE_NULL_HANDLING));
   }
diff --git 
a/pinot-common/src/test/java/org/apache/pinot/common/utils/config/QueryOptionsUtilsTest.java
 
b/pinot-common/src/test/java/org/apache/pinot/common/utils/config/QueryOptionsUtilsTest.java
index f2125adef28..ffb698e789d 100644
--- 
a/pinot-common/src/test/java/org/apache/pinot/common/utils/config/QueryOptionsUtilsTest.java
+++ 
b/pinot-common/src/test/java/org/apache/pinot/common/utils/config/QueryOptionsUtilsTest.java
@@ -39,7 +39,7 @@ import static org.testng.Assert.fail;
 public class QueryOptionsUtilsTest {
   private static final List<String> POSITIVE_INT_KEYS =
       List.of(NUM_REPLICA_GROUPS_TO_QUERY, MAX_EXECUTION_THREADS, 
NUM_GROUPS_LIMIT, MAX_INITIAL_RESULT_HOLDER_CAPACITY,
-          MAX_STREAMING_PENDING_BLOCKS, MAX_ROWS_IN_JOIN, MAX_ROWS_IN_WINDOW);
+          MAX_STREAMING_PENDING_BLOCKS, MAX_ROWS_IN_JOIN, MAX_ROWS_IN_WINDOW, 
SORTED_SELECTION_MERGE_BLOCK_SIZE);
   private static final List<String> NON_NEGATIVE_INT_KEYS =
       List.of(MULTI_STAGE_LEAF_LIMIT, STREAMING_GROUP_BY_FLUSH_THRESHOLD, 
STREAMING_DISTINCT_FLUSH_THRESHOLD);
   private static final List<String> UNBOUNDED_INT_KEYS =
@@ -375,6 +375,8 @@ public class QueryOptionsUtilsTest {
         return QueryOptionsUtils.getMaxRowsInJoin(map);
       case MAX_ROWS_IN_WINDOW:
         return QueryOptionsUtils.getMaxRowsInWindow(map);
+      case SORTED_SELECTION_MERGE_BLOCK_SIZE:
+        return QueryOptionsUtils.getSortedSelectionMergeBlockSize(map);
       // Non-negative ints
       case MULTI_STAGE_LEAF_LIMIT:
         return QueryOptionsUtils.getMultiStageLeafLimit(map);
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
new file mode 100644
index 00000000000..b99d6f0edc7
--- /dev/null
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/operator/combine/StreamingSelectionOrderByCombineOperator.java
@@ -0,0 +1,543 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.core.operator.combine;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Comparator;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.PriorityQueue;
+import java.util.Set;
+import java.util.concurrent.ExecutorService;
+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.core.common.Operator;
+import org.apache.pinot.core.operator.AcquireReleaseColumnsSegmentOperator;
+import org.apache.pinot.core.operator.blocks.results.BaseResultsBlock;
+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.operator.transform.function.TransformFunction;
+import 
org.apache.pinot.core.operator.transform.function.TransformFunctionFactory;
+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;
+import org.apache.pinot.segment.spi.IndexSegment;
+import org.apache.pinot.segment.spi.datasource.DataSource;
+import org.apache.pinot.segment.spi.datasource.DataSourceMetadata;
+import org.apache.pinot.spi.exception.QueryErrorCode;
+import org.apache.pinot.spi.exception.QueryErrorMessage;
+import org.apache.pinot.spi.query.QueryThreadContext;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/// Streaming, lazy combine operator for selection ORDER BY queries whose 
first order-by expression is an identifier.
+///
+/// It performs an incremental k-way heap merge across the per-segment 
operators in {@code _operators}, returning
+/// globally sorted rows in bounded blocks. Each segment is exposed through a 
{@link SegmentCursor} that yields that
+/// segment's locally-sorted rows in order:
+///
+/// - Segments physically sorted on the first order-by column are backed by
+///   {@link StreamingSelectionOrderByOperator}, which is pulled lazily one 
run/block at a time.
+/// - Other (e.g. consuming/unsorted) segments are backed by a single 
materialized top-K block (any
+///   {@link SelectionResultsBlock}-producing operator such as {@code 
SelectionOrderByOperator}); the cursor reads that
+///   one block and iterates its rows.
+///
+/// A {@link PriorityQueue} of {@link SegmentCursor} ordered by the {@link 
OrderByComparatorFactory} comparator on
+/// each
+/// cursor's current head row drives the merge with an 
at-most-one-head-per-active-segment invariant (the heap holds the
+/// cursors themselves, never all rows, which would degenerate into a full 
heap-sort that materializes everything). Each
+/// cycle pops the global-min cursor, appends its head to the current output 
block, advances that one cursor by a single
+/// row, and re-offers it if it still has a head.
+///
+/// **Min/max lazy segment activation (pruning).** Cursors are sorted by the 
first order-by column's min value
+/// (ASC) / max value (DESC) reusing the {@code MinMaxValueContext} idea from
+/// {@link MinMaxValueBasedSelectionOrderByCombineOperator}. A cursor is only 
activated (its segment acquired and first
+/// block read) when the merge frontier reaches its min/max, so once {@code 
limit + offset} rows are emitted the
+/// remaining segments are never acquired or read. See {@link 
#activateEligibleCursors()} 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.
+///
+/// **Segment acquire/release lifecycle.** A cursor acquires its
+/// {@link AcquireReleaseColumnsSegmentOperator} on activation and releases it 
only when its child operator is fully
+/// drained (acquire-on-activate / release-on-exhaust), rather than per run. 
This is intentional: the backing
+/// {@link StreamingSelectionOrderByOperator} retains a buffer-backed {@code 
ValueBlock} across {@code nextBlock()}
+/// calls in its tail-to-sort mode, so releasing between interleaved runs 
could read segment buffers after a release
+/// under prefetch. Holding the acquire for the cursor's lifetime guarantees 
no release happens between a cursor's
+/// own reads; min/max pruning bounds the number of simultaneously-active 
(acquired) segments to the merge frontier.
+/// The rows handed out by the child operators are already deep-copied to heap 
{@code Object[]} (via
+/// {@code RowBasedBlockValueFetcher}), so they remain valid after the segment 
is released. Any cursors still
+/// acquired when the merge ends early (LIMIT reached) or errors out are 
released via {@link #releaseAllCursors()}.
+///
+/// **Streaming vs single-block.** When {@code _streaming} is {@code true} 
(MSE leaf path, driven by
+/// {@link 
org.apache.pinot.core.operator.streaming.StreamingInstanceResponseOperator}) 
the merge emits many bounded
+/// {@link SelectionResultsBlock}s from successive {@link #getNextBlock()} 
calls followed by a final
+/// {@link MetadataResultsBlock}. When {@code false} (classic single-stage 
path) the merge runs to completion and the
+/// first {@link #getNextBlock()} call returns a single block with execution 
stats attached.
+///
+/// **Threading.** This operator overrides {@link #start()}/{@link #stop()} to 
no-ops (other than releasing
+/// segments) and runs the merge single-threaded and lazily in {@link 
#getNextBlock()} on the consumer thread; it does
+/// not use the base worker-queue model, and {@link #processSegments()} is 
overridden to fail loud. The base
+/// {@code Phaser} (which exists only to fence worker threads against segment 
release) is intentionally bypassed because
+/// all child/segment access is synchronous on the single consumer thread that 
holds the segment references; no async
+/// work may be introduced here without restoring that fence. The instance is 
single-use (driven once to completion) and
+/// is not thread-safe.
+@SuppressWarnings({"rawtypes", "unchecked"})
+public class StreamingSelectionOrderByCombineOperator extends 
BaseStreamingCombineOperator<SelectionResultsBlock> {
+  private static final Logger LOGGER = 
LoggerFactory.getLogger(StreamingSelectionOrderByCombineOperator.class);
+  private static final String EXPLAIN_NAME = 
"COMBINE_SELECT_ORDERBY_STREAMING";
+
+  private final boolean _streaming;
+  private final boolean _asc;
+  private final boolean _pruningEnabled;
+  private final int _numRowsToKeep;
+  private final int _blockSize;
+  private final Comparator<Object[]> _comparator;
+  private final SegmentCursor[] _sortedCursors;
+  private final PriorityQueue<SegmentCursor> _priorityQueue;
+
+  // Merge progress (single-threaded; mutated only by the consumer thread 
driving getNextBlock())
+  private int _nextToActivate;
+  private int _numRowsEmitted;
+  private List<Object[]> _outputRows;
+  private boolean _done;
+  /// Captured from the first child block seen; all child blocks share the 
same schema
+  private DataSchema _dataSchema;
+  /// Deduplicated MERGE_RESPONSE errors for segment blocks dropped on schema 
mismatch; null until the first mismatch
+  @Nullable
+  private Set<String> _dataSchemaMismatchErrors;
+  /// Subset of the above not yet attached to an emitted block; drained on 
each attach
+  @Nullable
+  private List<String> _unreportedDataSchemaMismatchErrors;
+
+  public StreamingSelectionOrderByCombineOperator(List<Operator> operators, 
QueryContext queryContext,
+      ExecutorService executorService, boolean streaming) {
+    // Pass a null merger: we override the consumption path entirely and never 
touch the base merger / worker queue.
+    super(null, operators, queryContext, executorService);
+    _streaming = streaming;
+    _numRowsToKeep = queryContext.getLimit() + queryContext.getOffset();
+    // Streaming mode flushes bounded blocks; single-stage mode flushes once 
at the end as a single block.
+    _blockSize = streaming ? queryContext.getSortedSelectionMergeBlockSize() : 
Integer.MAX_VALUE;
+    _pruningEnabled = !queryContext.isNullHandlingEnabled();
+
+    List<OrderByExpressionContext> orderByExpressions = 
queryContext.getOrderByExpressions();
+    assert orderByExpressions != null && !orderByExpressions.isEmpty();
+    OrderByExpressionContext firstOrderByExpression = 
orderByExpressions.get(0);
+    assert firstOrderByExpression.getExpression().getType() == 
ExpressionContext.Type.IDENTIFIER;
+    _asc = firstOrderByExpression.isAsc();
+    String firstOrderByColumn = 
firstOrderByExpression.getExpression().getIdentifier();
+    _comparator = OrderByComparatorFactory.getComparator(orderByExpressions, 
queryContext.isNullHandlingEnabled());
+
+    // Build one cursor per segment operator and read its first order-by 
column min/max for lazy activation ordering.
+    // Reading DataSourceMetadata does not touch column buffers, so no segment 
acquire is needed here (mirrors
+    // MinMaxValueBasedSelectionOrderByCombineOperator).
+    _sortedCursors = new SegmentCursor[_numOperators];
+    for (int i = 0; i < _numOperators; i++) {
+      Operator<BaseResultsBlock> operator = _operators.get(i);
+      DataSourceMetadata metadata =
+          operator.getIndexSegment().getDataSource(firstOrderByColumn, 
queryContext.getSchema())
+              .getDataSourceMetadata();
+      _sortedCursors[i] = new SegmentCursor(operator, metadata.getMinValue(), 
metadata.getMaxValue());
+    }
+    sortCursorsByMinMax();
+
+    _priorityQueue = new PriorityQueue<>(Math.max(1, _numOperators),
+        (o1, o2) -> _comparator.compare(o1.currentHead(), o2.currentHead()));
+    _outputRows = newOutputList();
+  }
+
+  /// Sorts the cursors so the merge can activate them lazily in frontier 
order: ascending by the column min value for
+  /// ASC, descending by the column max value for DESC. Cursors without a 
min/max are placed first because they must
+  /// always be processed (mirrors {@link 
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);
+      });
+    }
+  }
+
+  @Override
+  public String toExplainString() {
+    return EXPLAIN_NAME;
+  }
+
+  /// Override to a no-op: the merge is single-threaded and lazy in {@link 
#getNextBlock()}, so we do not spin up the
+  /// base worker threads / blocking-queue model.
+  @Override
+  public void start() {
+  }
+
+  /// Override the base worker-queue stop: no worker threads / phaser tasks 
were started. Release any segments still
+  /// acquired (idempotent) so an early stop by the driver cannot leak 
acquires.
+  @Override
+  public void stop() {
+    _done = true;
+    releaseAllCursors();
+  }
+
+  /// The base worker-thread entry point must never run here ({@link #start()} 
is a no-op). Fail loud if it ever does.
+  @Override
+  protected void processSegments() {
+    throw new IllegalStateException(
+        "StreamingSelectionOrderByCombineOperator runs single-threaded; 
processSegments() must not be called");
+  }
+
+  @Override
+  protected BaseResultsBlock getNextBlock() {
+    if (_done) {
+      // Streaming mode: terminal metadata block after the last data block. 
Idempotent if called again.
+      return attachExecutionStats(new MetadataResultsBlock());
+    }
+    try {
+      long endTimeMs = _queryContext.getEndTimeMs();
+      while (_numRowsEmitted < _numRowsToKeep) {
+        // The merge drains the heap on this thread and, for single-block 
cursors, may never re-enter a child operator.
+        // Without this check nothing would observe cancellation, pause, or 
the query deadline for up to
+        // (limit + offset) iterations -- a regression against 
MinMaxValueBasedSelectionOrderByCombineOperator, which
+        // this operator replaces on the same query shape.
+        
QueryThreadContext.checkTerminationAndSampleUsagePeriodically(_numRowsEmitted, 
EXPLAIN_NAME, endTimeMs);
+        activateEligibleCursors();
+        SegmentCursor cursor = _priorityQueue.poll();
+        if (cursor == null) {
+          // All active cursors exhausted (and pruning guarantees the rest 
cannot contribute).
+          break;
+        }
+        _outputRows.add(cursor.currentHead());
+        _numRowsEmitted++;
+        cursor.advance();
+        if (cursor.currentHead() != null) {
+          _priorityQueue.offer(cursor);
+        }
+        if (_streaming && _outputRows.size() >= _blockSize) {
+          return flushDataBlock();
+        }
+      }
+      // Merge complete.
+      finish();
+      if (!_outputRows.isEmpty()) {
+        // Streaming: the final partial data block (next call returns the 
terminal metadata block).
+        // Single-stage: the single complete block.
+        return flushDataBlock();
+      }
+      if (_streaming) {
+        return attachDataSchemaMismatchErrors(attachExecutionStats(new 
MetadataResultsBlock()));
+      }
+      // Single-stage with no rows: still return a (single) block carrying the 
schema and execution stats.
+      return attachDataSchemaMismatchErrors(attachExecutionStats(
+          new SelectionResultsBlock(resolveDataSchema(), List.of(), 
_comparator, _queryContext)));
+    } catch (Exception e) {
+      _done = true;
+      releaseAllCursors();
+      return createExceptionResultsBlockAndAttachExecutionStats(e, "merging 
sorted selection results");
+    } catch (Throwable t) {
+      // An Error (e.g. OutOfMemoryError while accumulating output rows) must 
not leave segments acquired for the
+      // lifetime of the server. In single-stage mode there is no 
start()/stop() backstop around this operator.
+      _done = true;
+      releaseAllCursors();
+      throw t;
+    }
+  }
+
+  /// Activates not-yet-active cursors whose min/max value can still 
contribute before the current merge frontier.
+  ///
+  /// Cursors are visited in min/max-sorted order and {@code _nextToActivate} 
is advanced only when a cursor is
+  /// actually activated; a {@code break} merely defers the current cursor, 
which is re-evaluated against the (rising)
+  /// frontier on every subsequent call. Correctness: a not-yet-activated 
cursor whose first order-by value range starts
+  /// strictly beyond the current heap head cannot contain a row that sorts 
before that head (the first order-by column
+  /// is the primary sort key), so the head is the true global minimum and is 
safe to emit; the deferred cursor is
+  /// activated later, exactly when the frontier reaches its min/max. When the 
heap is empty the frontier is unknown, so
+  /// activation is forced (never pruned), which also drains any segments that 
sort entirely after the ones seen so far.
+  private void activateEligibleCursors() {
+    while (_nextToActivate < _sortedCursors.length) {
+      SegmentCursor cursor = _sortedCursors[_nextToActivate];
+      if (_pruningEnabled) {
+        Comparable bound = _asc ? cursor._minValue : cursor._maxValue;
+        // A null bound means the segment must always be processed. Otherwise, 
only prune against a non-null frontier;
+        // if the head's first order-by value is null we cannot compare, so 
fall through and activate.
+        if (bound != null && !_priorityQueue.isEmpty()) {
+          Object headValue = _priorityQueue.peek().currentHead()[0];
+          if (headValue != null) {
+            // Both come from the same first order-by column: the metadata 
min/max and the materialized row[0] share
+            // the column's stored type, so this comparison is type-safe (same 
assumption as MinMaxValueBased...).
+            int cmp = bound.compareTo(headValue);
+            if (_asc ? cmp > 0 : cmp < 0) {
+              break;
+            }
+          }
+        }
+      }
+      cursor.activate();
+      _nextToActivate++;
+      if (cursor.currentHead() != null) {
+        _priorityQueue.offer(cursor);
+      }
+    }
+  }
+
+  /// Marks the merge done and releases any segments still held by un-drained 
cursors (e.g. when LIMIT is reached).
+  private void finish() {
+    _done = true;
+    releaseAllCursors();
+  }
+
+  private void releaseAllCursors() {
+    for (SegmentCursor cursor : _sortedCursors) {
+      cursor.release();
+    }
+  }
+
+  /// Returns the accumulated output rows as a sorted {@link 
SelectionResultsBlock} and resets the output buffer. The
+  /// block carries the comparator so the broker-side n-way reduce stays 
correct. In single-stage mode it is the only
+  /// block, so execution stats are attached; in streaming mode stats are 
attached to the terminal metadata block.
+  private BaseResultsBlock flushDataBlock() {
+    List<Object[]> rows = _outputRows;
+    _outputRows = newOutputList();
+    SelectionResultsBlock block = new 
SelectionResultsBlock(resolveDataSchema(), rows, _comparator, _queryContext);
+    attachDataSchemaMismatchErrors(block);
+    return _streaming ? block : attachExecutionStats(block);
+  }
+
+  /// Records that a segment's block was dropped because its schema disagreed 
with the merge schema. Deduplicated by
+  /// message so a mid-reload table with many divergent segments cannot flood 
the response.
+  private void recordDataSchemaMismatch(@Nullable DataSchema mismatched) {
+    String errorMessage =
+        String.format("Data schema mismatch between merged block: %s and block 
to merge: %s, drop block to merge",
+            _dataSchema, mismatched);
+    // NOTE: This is segment level log, so log at debug level to prevent 
flooding the log.
+    LOGGER.debug(errorMessage);
+    if (_dataSchemaMismatchErrors == null) {
+      _dataSchemaMismatchErrors = new HashSet<>();
+      _unreportedDataSchemaMismatchErrors = new ArrayList<>();
+    }
+    if (_dataSchemaMismatchErrors.add(errorMessage)) {
+      _unreportedDataSchemaMismatchErrors.add(errorMessage);
+    }
+  }
+
+  /// Attaches mismatch errors recorded since the last emitted block. In 
streaming mode blocks are emitted many times,
+  /// so errors are drained rather than re-attached, and the terminal block 
picks up any recorded after the last flush.
+  private <T extends BaseResultsBlock> T attachDataSchemaMismatchErrors(T 
block) {
+    if (_unreportedDataSchemaMismatchErrors != null && 
!_unreportedDataSchemaMismatchErrors.isEmpty()) {
+      for (String errorMessage : _unreportedDataSchemaMismatchErrors) {
+        
block.addErrorMessage(QueryErrorMessage.safeMsg(QueryErrorCode.MERGE_RESPONSE, 
errorMessage));
+      }
+      _unreportedDataSchemaMismatchErrors.clear();
+    }
+    return block;
+  }
+
+  private List<Object[]> newOutputList() {
+    int capacity =
+        Math.min(_blockSize, Math.min(_numRowsToKeep, 
SelectionOperatorUtils.MAX_ROW_HOLDER_INITIAL_CAPACITY));
+    return new ArrayList<>(Math.max(1, capacity));
+  }
+
+  /// Returns the data schema captured from the first child block. If no 
segment produced a block (every segment is a
+  /// streaming operator that matched zero rows), reconstructs the schema from 
the first segment so that an empty result
+  /// still carries a valid, correctly-ordered schema (order-by expressions 
first, matching the child blocks' layout).
+  private DataSchema resolveDataSchema() {
+    if (_dataSchema != null) {
+      return _dataSchema;
+    }
+    IndexSegment indexSegment = _operators.get(0).getIndexSegment();
+    List<ExpressionContext> expressions = 
SelectionOperatorUtils.extractExpressions(_queryContext, indexSegment);
+    Set<String> columns = new HashSet<>();
+    for (ExpressionContext expression : expressions) {
+      expression.getColumns(columns);
+    }
+    Map<String, DataSource> dataSourceMap = new HashMap<>();
+    for (String column : columns) {
+      dataSourceMap.put(column, indexSegment.getDataSource(column, 
_queryContext.getSchema()));
+    }
+    int numExpressions = expressions.size();
+    String[] columnNames = new String[numExpressions];
+    ColumnDataType[] columnDataTypes = new ColumnDataType[numExpressions];
+    for (int i = 0; i < numExpressions; i++) {
+      ExpressionContext expression = expressions.get(i);
+      columnNames[i] = expression.toString();
+      TransformFunction transformFunction = 
TransformFunctionFactory.get(expression, dataSourceMap);
+      columnDataTypes[i] = 
ColumnDataType.fromDataType(transformFunction.getResultMetadata().getDataType(),
+          transformFunction.getResultMetadata().isSingleValue());
+    }
+    _dataSchema = new DataSchema(columnNames, columnDataTypes);
+    return _dataSchema;
+  }
+
+  /// Iterates a single segment's locally-sorted rows, pulling blocks lazily 
from its operator. Streaming-backed cursors
+  /// loop until the operator returns {@code null}; single-block-backed 
cursors read exactly one block. The segment is
+  /// acquired on activation and released once exhausted or when the combine 
finishes (see
+  /// {@link StreamingSelectionOrderByCombineOperator}).
+  private class SegmentCursor {
+    private final Operator<BaseResultsBlock> _operator;
+    @Nullable
+    private final Comparable _minValue;
+    @Nullable
+    private final Comparable _maxValue;
+
+    private boolean _streamingChild;
+    private List<Object[]> _rows;
+    private int _pos;
+    private boolean _acquired;
+    /// Cached current head (rows.get(pos)) so the hot heap comparator does 
not re-index per comparison
+    @Nullable
+    private Object[] _head;
+
+    SegmentCursor(Operator<BaseResultsBlock> operator, @Nullable Comparable 
minValue, @Nullable Comparable maxValue) {
+      _operator = operator;
+      _minValue = minValue;
+      _maxValue = maxValue;
+    }
+
+    /// Returns the current head row to be merged next, or {@code null} if not 
activated or exhausted.
+    @Nullable
+    Object[] currentHead() {
+      return _head;
+    }
+
+    /// Acquires the segment, resolves whether the child is the lazy streaming 
operator, and reads the first block.
+    /// After this call {@link #currentHead()} returns the first row, or 
{@code null} if the segment contributes
+    /// nothing.
+    void activate() {
+      acquireSegment();
+      _streamingChild = isStreamingChild();
+      if (!pullBlock()) {
+        exhaust();
+      }
+    }
+
+    /// Advances past the current head, pulling the next block lazily for 
streaming cursors. Releases the segment
+    /// when the cursor is exhausted; afterwards {@link #currentHead()} 
returns {@code null}.
+    void advance() {
+      _pos++;
+      if (_pos < _rows.size()) {
+        _head = _rows.get(_pos);
+        return;
+      }
+      // Current block drained: streaming cursors pull the next run/block; 
single-block cursors are done.
+      if (_streamingChild && pullBlock()) {
+        return;
+      }
+      exhaust();
+    }
+
+    /// Loads the next non-empty block of rows from the operator, capturing 
the combine-level data schema on first
+    /// sight. Returns {@code false} when the operator is exhausted (no more 
rows). Streaming-backed operators emit
+    /// one run/block per call and {@code null} when done; single-block 
operators emit a single block and must not be
+    /// called again afterwards, so an empty/null block from a single-block 
child is treated as exhausted.
+    ///
+    /// A block whose schema differs from the one already captured is dropped 
and reported, mirroring
+    /// {@link 
org.apache.pinot.core.operator.combine.merger.SelectionOrderByResultsBlockMerger}.
 Segments on a server
+    /// can disagree on schema mid-reload (a newly added column exists only in 
reloaded segments), and merging rows of
+    /// differing width under one schema would corrupt the result rather than 
fail.
+    private boolean pullBlock() {
+      while (true) {
+        SelectionResultsBlock block = nextBlock();
+        if (block == null) {
+          return false;
+        }
+        if (_dataSchema == null) {
+          _dataSchema = block.getDataSchema();
+        } else if (!_dataSchema.equals(block.getDataSchema())) {
+          recordDataSchemaMismatch(block.getDataSchema());
+          return false;
+        }
+        List<Object[]> rows = block.getRows();
+        if (rows != null && !rows.isEmpty()) {
+          _rows = rows;
+          _pos = 0;
+          _head = rows.get(0);
+          return true;
+        }
+        // Defensive: an unexpected empty (non-null) block. Keep pulling only 
for streaming children; a single-block
+        // child yields exactly one block, so treat it as exhausted.
+        if (!_streamingChild) {
+          return false;
+        }
+      }
+    }
+
+    private SelectionResultsBlock nextBlock() {
+      try {
+        return (SelectionResultsBlock) _operator.nextBlock();
+      } catch (RuntimeException e) {
+        throw wrapOperatorException(_operator, e);
+      }
+    }
+
+    /// Returns whether the underlying child operator is the lazy {@link 
StreamingSelectionOrderByOperator}. Must be
+    /// called after {@link #acquireSegment()} because materializing the 
wrapped child runs the plan node, which
+    /// accesses segment buffers.
+    private boolean isStreamingChild() {
+      Operator underlying = _operator;
+      if (_operator instanceof AcquireReleaseColumnsSegmentOperator) {
+        AcquireReleaseColumnsSegmentOperator wrapper = 
(AcquireReleaseColumnsSegmentOperator) _operator;
+        wrapper.materializeChildOperator();
+        underlying = wrapper.getChildOperators().get(0);
+      }
+      return underlying instanceof StreamingSelectionOrderByOperator;
+    }
+
+    private void acquireSegment() {
+      if (_operator instanceof AcquireReleaseColumnsSegmentOperator) {
+        ((AcquireReleaseColumnsSegmentOperator) _operator).acquire();
+      }
+      _acquired = true;
+    }
+
+    /// Releases the segment if still held. Idempotent: safe to call from 
{@link #exhaust()} and combine cleanup.
+    private void release() {
+      if (_acquired) {
+        if (_operator instanceof AcquireReleaseColumnsSegmentOperator) {
+          ((AcquireReleaseColumnsSegmentOperator) _operator).release();
+        }
+        _acquired = false;
+      }
+    }
+
+    private void exhaust() {
+      release();
+      _rows = null;
+      _head = null;
+    }
+  }
+}
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
new file mode 100644
index 00000000000..e66a6ee0391
--- /dev/null
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/operator/query/StreamingSelectionOrderByOperator.java
@@ -0,0 +1,514 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.core.operator.query;
+
+import com.google.common.base.CaseFormat;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.Comparator;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.PriorityQueue;
+import java.util.Set;
+import java.util.stream.Collectors;
+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.core.common.BlockValSet;
+import org.apache.pinot.core.common.Operator;
+import org.apache.pinot.core.common.RowBasedBlockValueFetcher;
+import org.apache.pinot.core.operator.BaseOperator;
+import org.apache.pinot.core.operator.BaseProjectOperator;
+import org.apache.pinot.core.operator.BitmapDocIdSetOperator;
+import org.apache.pinot.core.operator.ColumnContext;
+import org.apache.pinot.core.operator.ExecutionStatistics;
+import org.apache.pinot.core.operator.ExplainAttributeBuilder;
+import org.apache.pinot.core.operator.ProjectionOperator;
+import org.apache.pinot.core.operator.ProjectionOperatorUtils;
+import org.apache.pinot.core.operator.blocks.ValueBlock;
+import org.apache.pinot.core.operator.blocks.results.SelectionResultsBlock;
+import org.apache.pinot.core.operator.transform.TransformOperator;
+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;
+import org.apache.pinot.segment.spi.IndexSegment;
+import org.apache.pinot.segment.spi.datasource.DataSource;
+import org.apache.pinot.spi.query.QueryScanCostContext;
+import org.roaringbitmap.RoaringBitmap;
+
+
+/// Lazy, incremental selection ORDER BY operator for segments that are 
physically sorted on the first order-by column.
+///
+/// Unlike {@link SelectionOrderByOperator} (which materializes the segment's 
whole top-K in a single block) this
+/// operator emits one globally-sorted {@link SelectionResultsBlock} per 
{@link #getNextBlock()} call and returns
+/// {@code null} when the segment is exhausted, so that a downstream 
k-way-merge combine operator can pull from many
+/// segments lazily and stop early. It relies on the underlying project 
operator iterating the first order-by column in
+/// the query order (the caller must guarantee {@code 
projectOperator.isCompatibleWith(DocIdOrder.fromAsc(asc))}).
+///
+/// It runs in one of two emission modes:
+///
+/// - **No tail to sort** ({@code numSortedExpressions == 
numOrderByExpressions}, e.g. {@code ORDER BY sorted}):
+///   rows already arrive from the project operator in final order, so each 
call emits the next project block
+///   (trimmed to the remaining {@code limit + offset} budget).
+/// - **Tail to sort** ({@code numSortedExpressions < numOrderByExpressions}, 
e.g.
+///   {@code ORDER BY sorted, other}):
+///   each call reads forward until the first order-by value changes (a 
primary-value "run"), retains the run's top
+///   {@code limit + offset} rows by the full comparator, and emits them 
sorted. This bounds the in-memory run buffer to
+///   {@code limit + offset} rows even when the first order-by column is 
near-constant (very low cardinality).
+///
+/// Like {@link SelectionOrderByOperator} it preserves the two-phase 
projection optimization: when there are output
+/// expressions that are not order-by expressions, the forward scan only 
fetches the order-by expressions plus the
+/// document id, and the non-order-by expressions are fetched in a second pass 
over the retained document ids of each
+/// emitted block.
+///
+/// This operator is stateful across {@link #getNextBlock()} calls and is 
**not** thread-safe; a single consumer
+/// must drive it.
+public class StreamingSelectionOrderByOperator extends 
BaseOperator<SelectionResultsBlock> {
+  private static final String EXPLAIN_NAME = "SELECT_ORDERBY_STREAMING";
+
+  private final IndexSegment _indexSegment;
+  private final QueryContext _queryContext;
+  private final boolean _nullHandlingEnabled;
+  /// Deduped order-by expressions followed by output expressions from 
SelectionOperatorUtils.extractExpressions()
+  private final List<ExpressionContext> _expressions;
+  private final BaseProjectOperator<?> _projectOperator;
+  private final List<OrderByExpressionContext> _orderByExpressions;
+  private final ColumnContext[] _orderByColumnContexts;
+  private final int _numExpressions;
+  private final int _numOrderByExpressions;
+  private final int _numRowsToKeep;
+  /// Whether there are output expressions that are not order-by expressions 
(requires the two-phase fetch)
+  private final boolean _twoPhase;
+  /// Whether the order-by has an unsorted tail that must be sorted in memory 
per run
+  private final boolean _tailToSort;
+  /// Expressions fetched during the forward scan: order-by expressions only 
when two-phase, otherwise all expressions
+  private final List<ExpressionContext> _phase1Expressions;
+  private final int _numPhase1Columns;
+  private final Comparator<Object[]> _comparator;
+  /// Compares only the first order-by column; used to detect primary-value 
run boundaries
+  private final Comparator<Object[]> _primaryComparator;
+  /// Pre-allocated run heap (cleared and reused each nextRun() call to avoid 
per-run allocation)
+  private final Comparator<Object[]> _reversedComparator;
+  private final PriorityQueue<Object[]> _runHeap;
+
+  // Pre-computed invariants for the two-phase fetch (null when single-phase)
+  private final List<ExpressionContext> _nonOrderByExpressions;
+  private final Map<String, DataSource> _phase2DataSourceMap;
+  private final int _phase2NumColumns;
+
+  /// Lazily built and cached; for two-phase it requires the transform 
operator's result column contexts
+  private DataSchema _dataSchema;
+
+  // Forward-scan cursor state (used by the tail-to-sort mode)
+  private ValueBlock _currentBlock;
+  private RowBasedBlockValueFetcher _currentFetcher;
+  private int[] _currentDocIds;
+  private RoaringBitmap[] _currentNullBitmaps;
+  private int _currentNumDocs;
+  private int _currentPos;
+  /// One-row lookahead: the first row of the next run, stashed when a run 
boundary is crossed
+  private Object[] _pendingRow;
+  private boolean _projectExhausted;
+
+  private boolean _exhausted;
+  private int _numRowsEmitted;
+  private int _numDocsScanned = 0;
+  private long _numEntriesScannedPostFilter = 0;
+
+  public StreamingSelectionOrderByOperator(IndexSegment indexSegment, 
QueryContext queryContext,
+      List<ExpressionContext> expressions, BaseProjectOperator<?> 
projectOperator, int numSortedExpressions) {
+    _indexSegment = indexSegment;
+    _queryContext = queryContext;
+    _nullHandlingEnabled = queryContext.isNullHandlingEnabled();
+    _expressions = expressions;
+    _projectOperator = projectOperator;
+
+    _orderByExpressions = queryContext.getOrderByExpressions();
+    assert _orderByExpressions != null;
+    _numExpressions = expressions.size();
+    _numOrderByExpressions = _orderByExpressions.size();
+    _orderByColumnContexts = new ColumnContext[_numOrderByExpressions];
+    for (int i = 0; i < _numOrderByExpressions; i++) {
+      ExpressionContext expression = 
_orderByExpressions.get(i).getExpression();
+      _orderByColumnContexts[i] = 
_projectOperator.getResultColumnContext(expression);
+    }
+
+    _numRowsToKeep = queryContext.getOffset() + queryContext.getLimit();
+    _twoPhase = _numExpressions > _numOrderByExpressions;
+    _tailToSort = numSortedExpressions < _numOrderByExpressions;
+    _comparator =
+        OrderByComparatorFactory.getComparator(_orderByExpressions, 
_orderByColumnContexts, _nullHandlingEnabled);
+    // The first order-by column is the physically sorted column, so it never 
contains nulls on this path; comparing
+    // only index 0 is enough to detect when one primary-value run ends and 
the next begins.
+    _primaryComparator =
+        OrderByComparatorFactory.getComparator(_orderByExpressions, 
_orderByColumnContexts, _nullHandlingEnabled, 0, 1);
+    _reversedComparator = _comparator.reversed();
+    _runHeap = new PriorityQueue<>(
+        Math.min(_numRowsToKeep, 
SelectionOperatorUtils.MAX_ROW_HOLDER_INITIAL_CAPACITY), _reversedComparator);
+
+    if (_twoPhase) {
+      _phase1Expressions = new ArrayList<>(_numOrderByExpressions);
+      for (OrderByExpressionContext orderByExpression : _orderByExpressions) {
+        _phase1Expressions.add(orderByExpression.getExpression());
+      }
+      _nonOrderByExpressions = _expressions.subList(_numOrderByExpressions, 
_numExpressions);
+      Set<String> columns = new HashSet<>();
+      for (ExpressionContext expressionContext : _nonOrderByExpressions) {
+        expressionContext.getColumns(columns);
+      }
+      _phase2NumColumns = columns.size();
+      _phase2DataSourceMap = new HashMap<>();
+      for (String column : columns) {
+        _phase2DataSourceMap.put(column, _indexSegment.getDataSource(column, 
_queryContext.getSchema()));
+      }
+    } else {
+      _phase1Expressions = _expressions;
+      _nonOrderByExpressions = null;
+      _phase2NumColumns = 0;
+      _phase2DataSourceMap = null;
+      // Single-phase: all output expressions are order-by expressions, so 
their types are known up front.
+      _dataSchema = buildSinglePhaseDataSchema();
+    }
+    _numPhase1Columns = _phase1Expressions.size();
+  }
+
+  @Override
+  protected SelectionResultsBlock getNextBlock() {
+    if (_exhausted) {
+      return null;
+    }
+    List<Object[]> rows = _tailToSort ? nextRun() : nextSortedRows();
+    if (rows == null || rows.isEmpty()) {
+      _exhausted = true;
+      return null;
+    }
+    if (_twoPhase) {
+      fetchNonOrderByColumns(rows);
+    }
+    // Single-phase builds the schema in the constructor; two-phase builds it 
during fetchNonOrderByColumns above.
+    assert _dataSchema != null;
+    return new SelectionResultsBlock(_dataSchema, rows, _comparator, 
_queryContext);
+  }
+
+  /// No-tail-to-sort mode: the project operator already returns rows in final 
order, so emit the next project block,
+  /// trimmed to the remaining {@code limit + offset} budget. Returns {@code 
null} when exhausted.
+  @Nullable
+  private List<Object[]> nextSortedRows() {
+    int remaining = _numRowsToKeep - _numRowsEmitted;
+    if (remaining <= 0) {
+      return null;
+    }
+    ValueBlock valueBlock = _projectOperator.nextBlock();
+    if (valueBlock == null) {
+      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);
+    List<Object[]> rows = new ArrayList<>(numRows);
+    for (int i = 0; i < numRows; i++) {
+      rows.add(materializeRow(blockValueFetcher, docIds, nullBitmaps, i));
+    }
+    _numRowsEmitted += rows.size();
+    return rows;
+  }
+
+  /// Tail-to-sort mode: read forward until the first order-by value changes, 
retain the run's top
+  /// {@code limit + offset} rows by the full comparator, and return them 
sorted. Returns {@code null} when
+  /// exhausted.
+  @Nullable
+  private List<Object[]> nextRun() {
+    int remaining = _numRowsToKeep - _numRowsEmitted;
+    if (remaining <= 0) {
+      return null;
+    }
+    if (_pendingRow == null) {
+      _pendingRow = nextRow();
+      if (_pendingRow == null) {
+        return null;
+      }
+    }
+    PriorityQueue<Object[]> runHeap = _runHeap;
+    runHeap.clear();
+    Object[] runFirstRow = _pendingRow;
+    SelectionOperatorUtils.addToPriorityQueue(_pendingRow, runHeap, 
_numRowsToKeep);
+    _pendingRow = null;
+    Object[] row;
+    while ((row = nextRow()) != null) {
+      if (_primaryComparator.compare(row, runFirstRow) == 0) {
+        SelectionOperatorUtils.addToPriorityQueue(row, runHeap, 
_numRowsToKeep);
+      } else {
+        // Run boundary: this row starts the next run, keep it for the next 
call.
+        _pendingRow = row;
+        break;
+      }
+    }
+    List<Object[]> rows = drainAscending(runHeap);
+    // A segment never contributes more than 'limit + offset' rows to the 
global result, and they are a prefix of its
+    // local sorted order, so cap the total emitted across runs at the 
remaining budget (the rows are ascending, keep
+    // the smallest 'remaining').
+    if (rows.size() > remaining) {
+      rows = rows.subList(0, remaining);
+    }
+    _numRowsEmitted += rows.size();
+    return rows;
+  }
+
+  /// Pulls the next row of the forward scan (across project blocks), 
materialized as an
+  /// {@code Object[_numExpressions]}. For two-phase the document id is 
stashed at index
+  /// {@code _numOrderByExpressions} (overwritten in the second pass). Returns 
{@code null} when the project operator
+  /// is exhausted.
+  @Nullable
+  private Object[] nextRow() {
+    while (true) {
+      if (_currentBlock == null || _currentPos >= _currentNumDocs) {
+        if (_projectExhausted) {
+          return null;
+        }
+        _currentBlock = _projectOperator.nextBlock();
+        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 (_currentNumDocs == 0) {
+          _currentBlock = null;
+          continue;
+        }
+      }
+      int rowId = _currentPos++;
+      return materializeRow(_currentFetcher, _currentDocIds, 
_currentNullBitmaps, rowId);
+    }
+  }
+
+  /// Materializes a single phase-1 row (deep-copied out of the value block 
buffers) from the given fetcher.
+  private Object[] materializeRow(RowBasedBlockValueFetcher fetcher, @Nullable 
int[] docIds,
+      @Nullable RoaringBitmap[] nullBitmaps, int rowId) {
+    Object[] row = new Object[_numExpressions];
+    fetcher.getRow(rowId, row, 0);
+    if (_twoPhase) {
+      row[_numOrderByExpressions] = docIds[rowId];
+    }
+    if (_nullHandlingEnabled) {
+      for (int colId = 0; colId < _numPhase1Columns; colId++) {
+        if (nullBitmaps[colId] != null && nullBitmaps[colId].contains(rowId)) {
+          row[colId] = null;
+        }
+      }
+    }
+    return row;
+  }
+
+  /// Drains a max-heap (created with the reversed comparator) into an 
ascending list, mutable so the second pass can
+  /// fill non-order-by values in place.
+  private List<Object[]> drainAscending(PriorityQueue<Object[]> heap) {
+    int numRows = heap.size();
+    Object[][] sortedRows = new Object[numRows][];
+    for (int i = numRows - 1; i >= 0; i--) {
+      sortedRows[i] = heap.poll();
+    }
+    return Arrays.asList(sortedRows);
+  }
+
+  /// Second pass of the two-phase fetch: fills the non-order-by expression 
values for the rows of a single emitted
+  /// block.
+  /// The rows keep their final (comparator) order; the fill iterates a 
document-id-sorted view that shares the same row
+  /// instances, mirroring {@link 
SelectionOrderByOperator#computePartiallyOrdered()}.
+  private void fetchNonOrderByColumns(List<Object[]> rows) {
+    int numRows = rows.size();
+    RoaringBitmap docIds = new RoaringBitmap();
+    for (Object[] row : rows) {
+      docIds.add((int) row[_numOrderByExpressions]);
+    }
+    // Document-id-sorted view sharing the same row instances (the bitmap 
returns docIds in ascending order).
+    List<Object[]> rowsByDocId = new ArrayList<>(rows);
+    rowsByDocId.sort(Comparator.comparingInt(o -> (int) 
o[_numOrderByExpressions]));
+
+    BitmapDocIdSetOperator docIdOperator = 
BitmapDocIdSetOperator.ascending(docIds, numRows);
+    try (ProjectionOperator projectionOperator =
+        ProjectionOperatorUtils.getProjectionOperator(_phase2DataSourceMap, 
docIdOperator, _queryContext)) {
+      TransformOperator transformOperator =
+          new TransformOperator(_queryContext, projectionOperator, 
_nonOrderByExpressions);
+
+      int numNonOrderByExpressions = _nonOrderByExpressions.size();
+      BlockValSet[] blockValSets = new BlockValSet[numNonOrderByExpressions];
+      int rowBaseId = 0;
+      ValueBlock valueBlock;
+      while ((valueBlock = transformOperator.nextBlock()) != null) {
+        for (int i = 0; i < numNonOrderByExpressions; i++) {
+          blockValSets[i] = 
valueBlock.getBlockValueSet(_nonOrderByExpressions.get(i));
+        }
+        RowBasedBlockValueFetcher blockValueFetcher = new 
RowBasedBlockValueFetcher(blockValSets);
+        int numDocsFetched = valueBlock.getNumDocs();
+        for (int i = 0; i < numDocsFetched; i++) {
+          blockValueFetcher.getRow(i, rowsByDocId.get(rowBaseId + i), 
_numOrderByExpressions);
+        }
+        if (_nullHandlingEnabled) {
+          RoaringBitmap[] nullBitmaps = new 
RoaringBitmap[numNonOrderByExpressions];
+          for (int i = 0; i < numNonOrderByExpressions; i++) {
+            nullBitmaps[i] = blockValSets[i].getNullBitmap();
+          }
+          for (int i = 0; i < numDocsFetched; i++) {
+            Object[] values = rowsByDocId.get(rowBaseId + i);
+            for (int colId = 0; colId < numNonOrderByExpressions; colId++) {
+              if (nullBitmaps[colId] != null && 
nullBitmaps[colId].contains(i)) {
+                values[_numOrderByExpressions + colId] = null;
+              }
+            }
+          }
+        }
+        _numEntriesScannedPostFilter += (long) numDocsFetched * 
_phase2NumColumns;
+        // Phase 2 re-reads docs already counted in phase 1, so only the extra 
entries are reported.
+        reportScanCost(0, (long) numDocsFetched * _phase2NumColumns);
+        rowBaseId += numDocsFetched;
+      }
+
+      if (_dataSchema == null) {
+        _dataSchema = buildTwoPhaseDataSchema(transformOperator);
+      }
+    }
+  }
+
+  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);
+  }
+
+  private DataSchema buildTwoPhaseDataSchema(TransformOperator 
transformOperator) {
+    int numNonOrderByExpressions = _nonOrderByExpressions.size();
+    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();
+    }
+    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 = 
transformOperator.getResultColumnContext(_nonOrderByExpressions.get(i));
+      columnDataTypes[_numOrderByExpressions + i] =
+          DataSchema.ColumnDataType.fromDataType(columnContext.getDataType(), 
columnContext.isSingleValue());
+    }
+    return new DataSchema(columnNames, columnDataTypes);
+  }
+
+  @Override
+  public String toExplainString() {
+    StringBuilder stringBuilder = new 
StringBuilder(EXPLAIN_NAME).append("(selectList:");
+    if (!_expressions.isEmpty()) {
+      stringBuilder.append(_expressions.get(0));
+      for (int i = 1; i < _expressions.size(); i++) {
+        stringBuilder.append(", ").append(_expressions.get(i));
+      }
+    }
+    return stringBuilder.append(')').toString();
+  }
+
+  @Override
+  protected String getExplainName() {
+    return CaseFormat.UPPER_UNDERSCORE.to(CaseFormat.UPPER_CAMEL, 
EXPLAIN_NAME);
+  }
+
+  @Override
+  protected void explainAttributes(ExplainAttributeBuilder attributeBuilder) {
+    super.explainAttributes(attributeBuilder);
+    if (_expressions.isEmpty()) {
+      return;
+    }
+    attributeBuilder.putStringList("selectList",
+        
_expressions.stream().map(ExpressionContext::toString).collect(Collectors.toList()));
+  }
+
+  @Override
+  public List<Operator> getChildOperators() {
+    return Collections.singletonList(_projectOperator);
+  }
+
+  @Override
+  public IndexSegment getIndexSegment() {
+    return _indexSegment;
+  }
+
+  @Override
+  public ExecutionStatistics getExecutionStatistics() {
+    long numEntriesScannedInFilter = 
_projectOperator.getExecutionStatistics().getNumEntriesScannedInFilter();
+    int numTotalDocs = _indexSegment.getSegmentMetadata().getTotalDocs();
+    return new ExecutionStatistics(_numDocsScanned, numEntriesScannedInFilter, 
_numEntriesScannedPostFilter,
+        numTotalDocs);
+  }
+
+  /// Reports scan cost to the shared {@link QueryScanCostContext} so 
scan-based query killing
+  /// ({@link org.apache.pinot.core.common.Operator#nextBlock()} -> {@code 
checkScanBasedKilling}) can see this
+  /// operator's work. Without this the killer is blind on the streaming path, 
unlike every other selection operator.
+  private void reportScanCost(int numDocsScanned, long 
numEntriesScannedPostFilter) {
+    QueryScanCostContext scanCost = getScanCostContext();
+    if (scanCost != null) {
+      if (numDocsScanned > 0) {
+        scanCost.addDocsScanned(numDocsScanned);
+      }
+      if (numEntriesScannedPostFilter > 0) {
+        scanCost.addEntriesScannedPostFilter(numEntriesScannedPostFilter);
+      }
+    }
+  }
+}
diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/plan/CombinePlanNode.java 
b/pinot-core/src/main/java/org/apache/pinot/core/plan/CombinePlanNode.java
index 71539035f53..298897e3ec2 100644
--- a/pinot-core/src/main/java/org/apache/pinot/core/plan/CombinePlanNode.java
+++ b/pinot-core/src/main/java/org/apache/pinot/core/plan/CombinePlanNode.java
@@ -34,6 +34,7 @@ import 
org.apache.pinot.core.operator.combine.SelectionOnlyCombineOperator;
 import org.apache.pinot.core.operator.combine.SelectionOrderByCombineOperator;
 import 
org.apache.pinot.core.operator.combine.SequentialSortedGroupByCombineOperator;
 import org.apache.pinot.core.operator.combine.SortedGroupByCombineOperator;
+import 
org.apache.pinot.core.operator.combine.StreamingSelectionOrderByCombineOperator;
 import 
org.apache.pinot.core.operator.streaming.StreamingDistinctCombineOperator;
 import 
org.apache.pinot.core.operator.streaming.StreamingGroupByCombineOperator;
 import 
org.apache.pinot.core.operator.streaming.StreamingSelectionOnlyCombineOperator;
@@ -131,6 +132,17 @@ public class CombinePlanNode implements PlanNode {
         // Use streaming operator only for non-empty selection-only query
         return createStreamingSelectionOnlyCombineOperator(operators);
       }
+      // Streaming selection order-by (opt-in via the 
sortedSelectionMergeEnabled hint). Selection-only already
+      // returned above, so reaching here with a non-empty limit and an 
order-by present implies selection order-by.
+      if (_queryContext.isSortedSelectionMergeEnabled() && 
QueryContextUtils.isSelectionQuery(_queryContext)
+          && _queryContext.getLimit() != 0) {
+        List<OrderByExpressionContext> orderByExpressions = 
_queryContext.getOrderByExpressions();
+        if (orderByExpressions != null
+            && orderByExpressions.get(0).getExpression().getType() == 
ExpressionContext.Type.IDENTIFIER) {
+          return new StreamingSelectionOrderByCombineOperator(operators, 
_queryContext, _executorService,
+              true /* streaming */);
+        }
+      }
       // Streaming flushes partial results, so it needs an aggregation above 
to merge them back together.
       // Leaves that must return final results are excluded, see 
StreamingGroupByCombineOperator.
       boolean leafReturnsFinalResult =
@@ -186,6 +198,10 @@ public class CombinePlanNode implements PlanNode {
     List<OrderByExpressionContext> orderByExpressions = 
_queryContext.getOrderByExpressions();
     assert orderByExpressions != null;
     if (orderByExpressions.get(0).getExpression().getType() == 
ExpressionContext.Type.IDENTIFIER) {
+      if (_queryContext.isSortedSelectionMergeEnabled()) {
+        return new StreamingSelectionOrderByCombineOperator(operators, 
_queryContext, _executorService,
+            false /* streaming */);
+      }
       return new MinMaxValueBasedSelectionOrderByCombineOperator(operators, 
_queryContext, _executorService);
     }
     return new SelectionOrderByCombineOperator(operators, _queryContext, 
_executorService);
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 029f4811c44..4849a193eeb 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
@@ -32,6 +32,7 @@ import 
org.apache.pinot.core.operator.query.SelectionOnlyOperator;
 import org.apache.pinot.core.operator.query.SelectionOrderByOperator;
 import 
org.apache.pinot.core.operator.query.SelectionPartiallyOrderedByDescOperation;
 import 
org.apache.pinot.core.operator.query.SelectionPartiallyOrderedByLinearOperator;
+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.IndexSegment;
@@ -88,11 +89,36 @@ public class SelectionPlanNode implements PlanNode {
         maxDocsPerCall = Math.min(limit + _queryContext.getOffset(), 
DocIdSetPlanNode.MAX_DOC_PER_CALL);
       }
 
-      BaseProjectOperator<?> projectOperator = getSortedByProject(expressions, 
maxDocsPerCall, orderByExpressions);
       boolean asc = orderByExpressions.get(0).isAsc();
       // Remember that we cannot use asc == projectOperator.isAscending() 
because empty operators are considered
       // both ascending and descending
       DocIdOrderedOperator.DocIdOrder queryOrder = 
DocIdOrderedOperator.DocIdOrder.fromAsc(asc);
+
+      // Opt-in streaming path: emit one globally-sorted block at a time so a 
downstream k-way-merge combine can pull
+      // lazily. Only build it when the first order-by column is an identifier 
(kept consistent with the combine-side
+      // gate) and the forward-scan project is order-compatible; the 
DESC-incompatible sorted case still falls back to
+      // the materialized SelectionPartiallyOrderedByDescOperation below so 
global order stays correct.
+      if (_queryContext.isSortedSelectionMergeEnabled()
+          && orderByExpressions.get(0).getExpression().getType() == 
ExpressionContext.Type.IDENTIFIER) {
+        // When there are non-order-by output expressions, only fetch the 
order-by expressions during the forward scan
+        // (the streaming operator fetches the rest in a second pass); 
otherwise fetch all expressions.
+        List<ExpressionContext> projectExpressions = expressions;
+        if (expressions.size() > numOrderByExpressions) {
+          projectExpressions = new ArrayList<>(numOrderByExpressions);
+          for (OrderByExpressionContext orderByExpression : 
orderByExpressions) {
+            projectExpressions.add(orderByExpression.getExpression());
+          }
+        }
+        BaseProjectOperator<?> streamingProjectOperator =
+            getSortedByProject(projectExpressions, maxDocsPerCall, 
orderByExpressions);
+        if (streamingProjectOperator.isCompatibleWith(queryOrder)) {
+          return new StreamingSelectionOrderByOperator(_indexSegment, 
_queryContext, expressions,
+              streamingProjectOperator, sortedColumnsPrefixSize);
+        }
+        // DESC-incompatible: fall through to the materialized fallback 
(rebuilds the project over all expressions).
+      }
+
+      BaseProjectOperator<?> projectOperator = getSortedByProject(expressions, 
maxDocsPerCall, orderByExpressions);
       if (projectOperator.isCompatibleWith(queryOrder)) {
         return new SelectionPartiallyOrderedByLinearOperator(_indexSegment, 
_queryContext, expressions, projectOperator,
             sortedColumnsPrefixSize);
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 3a24e0ed51e..25ff2a4db17 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
@@ -270,6 +270,17 @@ public class InstancePlanMakerImplV2 implements PlanMaker {
     }
     queryContext.setMaxExecutionThreads(maxExecutionThreads);
 
+    // Set streaming selection order-by options (opt-in; gated on selection 
queries to prevent accidental routing
+    // if a downstream guard is ever missed)
+    if (QueryContextUtils.isSelectionQuery(queryContext)) {
+      
queryContext.setSortedSelectionMergeEnabled(QueryOptionsUtils.isSortedSelectionMergeEnabled(queryOptions));
+      Integer sortedSelectionMergeBlockSize =
+          QueryOptionsUtils.getSortedSelectionMergeBlockSize(queryOptions);
+      if (sortedSelectionMergeBlockSize != null) {
+        
queryContext.setSortedSelectionMergeBlockSize(sortedSelectionMergeBlockSize);
+      }
+    }
+
     // Set group-by query options
     if (QueryContextUtils.isAggregationQuery(queryContext) && 
queryContext.getGroupByExpressions() != null) {
       // Set maxInitialResultHolderCapacity
diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/query/request/context/QueryContext.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/query/request/context/QueryContext.java
index 1371fd3f3c7..0d7f17775ce 100644
--- 
a/pinot-core/src/main/java/org/apache/pinot/core/query/request/context/QueryContext.java
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/query/request/context/QueryContext.java
@@ -44,6 +44,7 @@ import org.apache.pinot.core.util.MemoizedClassAssociation;
 import org.apache.pinot.segment.spi.datasource.DataSource;
 import org.apache.pinot.spi.config.table.FieldConfig;
 import org.apache.pinot.spi.data.Schema;
+import org.apache.pinot.spi.utils.CommonConstants.Broker;
 import org.apache.pinot.spi.utils.CommonConstants.Server;
 
 
@@ -135,6 +136,11 @@ public class QueryContext {
   private int _streamingGroupByFlushThreshold;
   // Flush threshold for streaming distinct (0 = disabled)
   private int _streamingDistinctFlushThreshold;
+  /// Opt-in: use the streaming k-way-merge selection ORDER BY combine over 
sorted segments
+  private boolean _sortedSelectionMergeEnabled;
+  /// Output block size (rows) for the streaming selection ORDER BY combine
+  private int _sortedSelectionMergeBlockSize = 
Broker.DEFAULT_SORTED_SELECTION_MERGE_BLOCK_SIZE;
+
   // Whether null handling is enabled
   private boolean _nullHandlingEnabled;
   // Whether server returns the final result
@@ -509,6 +515,22 @@ public class QueryContext {
     _streamingDistinctFlushThreshold = streamingDistinctFlushThreshold;
   }
 
+  public boolean isSortedSelectionMergeEnabled() {
+    return _sortedSelectionMergeEnabled;
+  }
+
+  public void setSortedSelectionMergeEnabled(boolean 
sortedSelectionMergeEnabled) {
+    _sortedSelectionMergeEnabled = sortedSelectionMergeEnabled;
+  }
+
+  public int getSortedSelectionMergeBlockSize() {
+    return _sortedSelectionMergeBlockSize;
+  }
+
+  public void setSortedSelectionMergeBlockSize(int 
sortedSelectionMergeBlockSize) {
+    _sortedSelectionMergeBlockSize = sortedSelectionMergeBlockSize;
+  }
+
   public boolean isNullHandlingEnabled() {
     return _nullHandlingEnabled;
   }
diff --git 
a/pinot-core/src/test/java/org/apache/pinot/core/operator/combine/CombineSlowOperatorsTest.java
 
b/pinot-core/src/test/java/org/apache/pinot/core/operator/combine/CombineSlowOperatorsTest.java
index 32a2e9c348b..115ea33a0e5 100644
--- 
a/pinot-core/src/test/java/org/apache/pinot/core/operator/combine/CombineSlowOperatorsTest.java
+++ 
b/pinot-core/src/test/java/org/apache/pinot/core/operator/combine/CombineSlowOperatorsTest.java
@@ -24,6 +24,7 @@ import java.util.concurrent.CountDownLatch;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicReference;
 import java.util.function.Supplier;
@@ -41,6 +42,7 @@ import 
org.apache.pinot.segment.spi.datasource.DataSourceMetadata;
 import org.apache.pinot.spi.exception.EarlyTerminationException;
 import org.apache.pinot.spi.exception.QueryErrorCode;
 import org.apache.pinot.spi.exception.QueryErrorMessage;
+import org.apache.pinot.spi.query.QueryThreadContext;
 import org.apache.pinot.util.TestUtils;
 import org.testng.Assert;
 import org.testng.annotations.AfterClass;
@@ -54,6 +56,7 @@ import static org.mockito.Mockito.when;
 import static org.testng.Assert.assertEquals;
 import static org.testng.Assert.assertFalse;
 import static org.testng.Assert.assertNotNull;
+import static org.testng.Assert.assertTrue;
 
 
 /// This test mimic the behavior of combining slow operators, where operation 
is not done by the timeout. When the
@@ -131,16 +134,7 @@ public class CombineSlowOperatorsTest {
   @Test
   public void testCancelMinMaxValueBasedSelectionOrderByCombineOperator() {
     CountDownLatch ready = new CountDownLatch(1);
-    List<Operator> operators = getOperators(ready, () -> {
-      IndexSegment seg = mock(IndexSegment.class);
-      DataSource ds = mock(DataSource.class);
-      DataSourceMetadata dsmd = mock(DataSourceMetadata.class);
-      when(dsmd.getMinValue()).thenReturn(100L);
-      when(dsmd.getMaxValue()).thenReturn(200L);
-      when(seg.getDataSource(anyString(), any())).thenReturn(ds);
-      when(ds.getDataSourceMetadata()).thenReturn(dsmd);
-      return seg;
-    });
+    List<Operator> operators = getOperators(ready, minMaxSegmentSupplier());
     QueryContext queryContext = 
QueryContextConverterUtils.getQueryContext("SELECT * FROM testTable ORDER BY 
column");
     queryContext.setEndTimeMs(System.currentTimeMillis() + 10000);
     MinMaxValueBasedSelectionOrderByCombineOperator combineOperator =
@@ -170,13 +164,85 @@ public class CombineSlowOperatorsTest {
     testCancelCombineOperator(combineOperator, ready);
   }
 
+  @Test
+  public void testCancelStreamingSelectionOrderByCombineOperator() {
+    CountDownLatch ready = new CountDownLatch(1);
+    List<Operator> operators = getOperators(ready, minMaxSegmentSupplier());
+    QueryContext queryContext = 
QueryContextConverterUtils.getQueryContext("SELECT * FROM testTable ORDER BY 
column");
+    queryContext.setEndTimeMs(System.currentTimeMillis() + 10000);
+    // Single-stage mode: nextBlock() drives the whole merge synchronously on 
the (cancellable) caller thread.
+    StreamingSelectionOrderByCombineOperator combineOperator =
+        new StreamingSelectionOrderByCombineOperator(operators, queryContext, 
_executorService, false);
+    testCancelCombineOperator(combineOperator, ready, operators);
+  }
+
+  @Test
+  public void 
testCancelStreamingSelectionOrderByCombineOperatorStreamingMode() {
+    CountDownLatch ready = new CountDownLatch(1);
+    List<Operator> operators = getOperators(ready, minMaxSegmentSupplier());
+    QueryContext queryContext = 
QueryContextConverterUtils.getQueryContext("SELECT * FROM testTable ORDER BY 
column");
+    queryContext.setEndTimeMs(System.currentTimeMillis() + 10000);
+    // Streaming (MSE-leaf) mode: the first getNextBlock() still drives the 
merge on the caller thread, so the same
+    // interrupt-on-cancel path applies and must surface an 
ExceptionResultsBlock.
+    StreamingSelectionOrderByCombineOperator combineOperator =
+        new StreamingSelectionOrderByCombineOperator(operators, queryContext, 
_executorService, true);
+    testCancelCombineOperator(combineOperator, ready, operators);
+  }
+
+  /// The merge loop drains its heap on the caller thread and, for 
single-block cursors, may never re-enter a child
+  /// operator, so it must check the query deadline itself. With an 
already-expired deadline the operator must surface a
+  /// timeout <b>before</b> activating any child - asserted via {@code 
_operationInProgress}, which also keeps the test
+  /// from passing vacuously if the timeout came from somewhere else.
+  @Test
+  public void testStreamingSelectionOrderByCombineOperatorHonorsDeadline() {
+    List<Operator> operators = getOperators(null, minMaxSegmentSupplier());
+    QueryContext queryContext = 
QueryContextConverterUtils.getQueryContext("SELECT * FROM testTable ORDER BY 
column");
+    queryContext.setEndTimeMs(System.currentTimeMillis() - 1);
+    StreamingSelectionOrderByCombineOperator combineOperator =
+        new StreamingSelectionOrderByCombineOperator(operators, queryContext, 
_executorService, false);
+    try (QueryThreadContext ignore = QueryThreadContext.openForSseTest()) {
+      BaseResultsBlock resultsBlock = combineOperator.nextBlock();
+      assertTrue(resultsBlock instanceof ExceptionResultsBlock,
+          "Expired deadline must surface as an ExceptionResultsBlock, got: " + 
resultsBlock.getClass().getName());
+    }
+    for (Operator operator : operators) {
+      assertFalse(((SlowOperator) operator)._operationInProgress.get(),
+          "Deadline must be observed before any child operator is driven");
+    }
+  }
+
+  /// A segment whose first order-by column reports a non-null min/max, as 
required by the min/max-based operators.
+  private static Supplier<IndexSegment> minMaxSegmentSupplier() {
+    return () -> {
+      IndexSegment seg = mock(IndexSegment.class);
+      DataSource ds = mock(DataSource.class);
+      DataSourceMetadata dsmd = mock(DataSourceMetadata.class);
+      when(dsmd.getMinValue()).thenReturn(100L);
+      when(dsmd.getMaxValue()).thenReturn(200L);
+      when(seg.getDataSource(anyString(), any())).thenReturn(ds);
+      when(ds.getDataSourceMetadata()).thenReturn(dsmd);
+      return seg;
+    };
+  }
+
   private void testCancelCombineOperator(BaseCombineOperator<?> 
combineOperator, CountDownLatch ready) {
+    testCancelCombineOperator(combineOperator, ready, List.of());
+  }
+
+  /// Submits the combine operator on a separate thread, waits for a child 
operator to start, cancels the future, and
+  /// asserts the operator surfaces an {@link ExceptionResultsBlock}. When 
{@code operatorsToVerify} is non-empty
+  /// (the single-threaded streaming combine, whose merge runs the children on 
the caller thread), it additionally
+  /// asserts the cancellation was genuine - a child actually started and none 
completed normally - so the test
+  /// cannot pass vacuously on an unrelated failure or by never exercising the 
cancel path.
+  private void testCancelCombineOperator(BaseCombineOperator<?> 
combineOperator, CountDownLatch ready,
+      List<Operator> operatorsToVerify) {
     AtomicReference<BaseResultsBlock> resultsBlock = new AtomicReference<>();
     // Avoid early finalization by not using Executors.newSingleThreadExecutor 
(java <= 20, JDK-8145304)
     ExecutorService combineExecutor = Executors.newFixedThreadPool(1);
     try {
       Future<?> future = combineExecutor.submit(() -> 
resultsBlock.set(combineOperator.nextBlock()));
-      ready.await();
+      // Bound the wait so a regression where no child operator ever starts 
fails fast here instead of hanging.
+      assertTrue(ready.await(10, TimeUnit.SECONDS), "Expected a child operator 
to start before cancellation");
       // At this point, the combineOperator is or will be waiting on 
future.get() for all sub operators, and the
       // waiting can be cancelled as below.
       future.cancel(true);
@@ -187,6 +253,12 @@ public class CombineSlowOperatorsTest {
     }
     TestUtils.waitForCondition((aVoid) -> resultsBlock.get() instanceof 
ExceptionResultsBlock, 10_000,
         "Should have been cancelled");
+    // Genuine-cancellation check: no child ran to normal completion (each was 
interrupted or never started), so the
+    // ExceptionResultsBlock above is the cancellation outcome rather than an 
unrelated error.
+    for (Operator operator : operatorsToVerify) {
+      assertFalse(((SlowOperator) operator)._notInterrupted.get(),
+          "No operator should have completed normally after cancellation");
+    }
   }
 
   /// NOTE: It is hard to test the logger behavior, but only one error message 
about the query timeout should be logged
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
new file mode 100644
index 00000000000..e2d9c9f8a9b
--- /dev/null
+++ 
b/pinot-core/src/test/java/org/apache/pinot/core/operator/combine/StreamingSelectionOrderByCombineOperatorTest.java
@@ -0,0 +1,553 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.core.operator.combine;
+
+import java.io.File;
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Comparator;
+import java.util.List;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.stream.Collectors;
+import org.apache.commons.io.FileUtils;
+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.core.common.Operator;
+import org.apache.pinot.core.operator.blocks.results.BaseResultsBlock;
+import org.apache.pinot.core.operator.blocks.results.MetadataResultsBlock;
+import org.apache.pinot.core.operator.blocks.results.SelectionResultsBlock;
+import org.apache.pinot.core.plan.CombinePlanNode;
+import org.apache.pinot.core.plan.PlanNode;
+import org.apache.pinot.core.plan.maker.InstancePlanMakerImplV2;
+import org.apache.pinot.core.plan.maker.PlanMaker;
+import org.apache.pinot.core.query.executor.ResultsBlockStreamer;
+import org.apache.pinot.core.query.request.context.QueryContext;
+import 
org.apache.pinot.core.query.request.context.utils.QueryContextConverterUtils;
+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.segment.creator.impl.SegmentIndexCreationDriverImpl;
+import org.apache.pinot.segment.local.segment.readers.GenericRowRecordReader;
+import org.apache.pinot.segment.spi.IndexSegment;
+import org.apache.pinot.segment.spi.SegmentContext;
+import org.apache.pinot.segment.spi.creator.SegmentGeneratorConfig;
+import org.apache.pinot.spi.config.table.TableConfig;
+import org.apache.pinot.spi.config.table.TableType;
+import org.apache.pinot.spi.data.FieldSpec;
+import org.apache.pinot.spi.data.Schema;
+import org.apache.pinot.spi.data.readers.GenericRow;
+import org.apache.pinot.spi.utils.CommonConstants.Server;
+import org.apache.pinot.spi.utils.ReadMode;
+import org.apache.pinot.spi.utils.builder.TableConfigBuilder;
+import org.intellij.lang.annotations.Language;
+import org.testng.annotations.AfterClass;
+import org.testng.annotations.BeforeClass;
+import org.testng.annotations.Test;
+
+import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertNotNull;
+import static org.testng.Assert.assertTrue;
+
+
+/// Combine-level tests for {@link StreamingSelectionOrderByCombineOperator} 
(step-3 operator) and its wiring into
+/// {@link CombinePlanNode#getCombineOperator()} (step-4).
+///
+/// <p>The streaming combine must return the same globally-sorted top-K rows 
as the default
+/// {@link MinMaxValueBasedSelectionOrderByCombineOperator}, only (in 
streaming mode) spread across several bounded
+/// blocks. Each functional test therefore asserts 
<b>streaming-vs-non-streaming parity</b>: it runs the identical query
+/// twice over the same in-memory segments - once with {@code 
sortedSelectionMergeEnabled=true} (asserting the new
+/// operator was actually selected) and once with the hint off (asserting the 
{@code MinMax} operator was selected) -
+/// then checks the two row sets are equal as a multiset and that the 
streaming output is fully sorted by the order-by
+/// comparator.
+///
+/// <p>To keep the top-K boundary unambiguous (operators may legitimately 
disagree on which of several rows that tie on
+/// every order-by key fall inside the limit) every parity query ends its 
ORDER BY with the globally-unique
+/// {@code valCol}
+/// so the comparator is a total order; the merge still genuinely interleaves 
segments because the primary sort column
+/// ({@code sortedCol}) overlaps across segments. Multiset (rather than 
positional) comparison then tolerates only the
+/// harmless reordering of fully-equal projected rows.
+public class StreamingSelectionOrderByCombineOperatorTest {
+  private static final File TEMP_DIR =
+      new File(FileUtils.getTempDirectory(), 
"StreamingSelectionOrderByCombineOperatorTest");
+  private static final String RAW_TABLE_NAME = "testTable";
+
+  private static final String SORTED_COL = "sortedCol";
+  private static final String TAIL_COL = "tailCol";
+  private static final String VAL_COL = "valCol";
+  private static final String NULLABLE_COL = "nullableCol";
+  /// Non-INT projected columns so the parity assertion can catch a 
stored-type / boxing regression (e.g. LONG emitted
+  /// where MinMax emits INT), which an all-INT suite cannot observe.
+  private static final String LONG_COL = "longCol";
+  private static final String STR_COL = "strCol";
+
+  /// Create (MAX_NUM_THREADS_PER_QUERY * 2) sorted segments so the leaf runs 
plan nodes across multiple threads.
+  private static final int NUM_SEGMENTS = 
QueryMultiThreadingUtils.MAX_NUM_THREADS_PER_QUERY * 2;
+  private static final int NUM_RECORDS_PER_SEGMENT = 100;
+
+  private static final TableConfig SORTED_TABLE_CONFIG =
+      new 
TableConfigBuilder(TableType.OFFLINE).setTableName(RAW_TABLE_NAME).setSortedColumn(SORTED_COL).build();
+  private static final TableConfig UNSORTED_TABLE_CONFIG =
+      new 
TableConfigBuilder(TableType.OFFLINE).setTableName(RAW_TABLE_NAME).build();
+  private static final Schema SCHEMA = new Schema.SchemaBuilder()
+      .addSingleValueDimension(SORTED_COL, FieldSpec.DataType.INT)
+      .addSingleValueDimension(TAIL_COL, FieldSpec.DataType.INT)
+      .addSingleValueDimension(VAL_COL, FieldSpec.DataType.INT)
+      .addSingleValueDimension(NULLABLE_COL, FieldSpec.DataType.INT)
+      .addSingleValueDimension(LONG_COL, FieldSpec.DataType.LONG)
+      .addSingleValueDimension(STR_COL, FieldSpec.DataType.STRING)
+      .build();
+
+  private static final PlanMaker PLAN_MAKER = new InstancePlanMakerImplV2();
+  private static final ExecutorService EXECUTOR = 
Executors.newCachedThreadPool();
+
+  /// Sorted segments with overlapping primary-column ranges, so the k-way 
merge interleaves them (genuine merge rather
+  /// than concatenation). Built with null handling on so the null-handling 
test sees real nulls in NULLABLE_COL; reads
+  /// with null handling off fall back to the column default.
+  private List<IndexSegment> _sortedSegments;
+  /// Sorted segments with disjoint, globally-increasing ranges: ORDER BY 
sortedCol with a small LIMIT drains only the
+  /// lowest segment, so min/max pruning must skip the rest (none 
acquired/scanned).
+  private List<IndexSegment> _disjointSegments;
+  /// A mix of sorted (streaming child) and physically-unsorted (single 
materialized top-K block child) segments,
+  /// exercising both SegmentCursor backings in one merge.
+  private List<IndexSegment> _mixedSegments;
+  /// Very low cardinality primary column (4 distinct values across 100 rows) 
so each value is a long run: exercises the
+  /// run/heap path in the streaming children and ties on the primary key at 
the prune boundary.
+  private List<IndexSegment> _lowCardSegments;
+
+  @BeforeClass
+  public void setUp()
+      throws Exception {
+    FileUtils.deleteDirectory(TEMP_DIR);
+
+    _sortedSegments = new ArrayList<>(NUM_SEGMENTS);
+    for (int i = 0; i < NUM_SEGMENTS; i++) {
+      _sortedSegments.add(buildSegment(SORTED_TABLE_CONFIG, "sorted_" + i, 
buildOverlappingSortedRecords(i), true));
+    }
+
+    _disjointSegments = new ArrayList<>(NUM_SEGMENTS);
+    for (int i = 0; i < NUM_SEGMENTS; i++) {
+      _disjointSegments.add(buildSegment(SORTED_TABLE_CONFIG, "disjoint_" + i, 
buildDisjointSortedRecords(i), false));
+    }
+
+    _lowCardSegments = new ArrayList<>(NUM_SEGMENTS);
+    for (int i = 0; i < NUM_SEGMENTS; i++) {
+      _lowCardSegments.add(buildSegment(SORTED_TABLE_CONFIG, "lowCard_" + i, 
buildLowCardinalityRecords(i), false));
+    }
+
+    // Two sorted + two unsorted segments, globally-unique valCol across all 
four so the multiset comparison is exact.
+    _mixedSegments = new ArrayList<>(4);
+    _mixedSegments.add(buildSegment(SORTED_TABLE_CONFIG, "mixedSorted_0", 
buildOverlappingSortedRecords(0), false));
+    _mixedSegments.add(buildSegment(SORTED_TABLE_CONFIG, "mixedSorted_1", 
buildOverlappingSortedRecords(1), false));
+    _mixedSegments.add(buildSegment(UNSORTED_TABLE_CONFIG, "mixedUnsorted_0", 
buildUnsortedRecords(2), false));
+    _mixedSegments.add(buildSegment(UNSORTED_TABLE_CONFIG, "mixedUnsorted_1", 
buildUnsortedRecords(3), false));
+  }
+
+  private static List<GenericRow> buildOverlappingSortedRecords(int index) {
+    int baseValue = index * NUM_RECORDS_PER_SEGMENT / 2;
+    List<GenericRow> records = new ArrayList<>(NUM_RECORDS_PER_SEGMENT);
+    for (int i = 0; i < NUM_RECORDS_PER_SEGMENT; i++) {
+      GenericRow record = new GenericRow();
+      record.putValue(SORTED_COL, baseValue + i);
+      record.putValue(TAIL_COL, NUM_RECORDS_PER_SEGMENT - i);
+      // Globally unique across all segments -> a total order when used as the 
final order-by key.
+      record.putValue(VAL_COL, index * 1_000_000 + i);
+      // Beyond the int range so a regression that narrows LONG -> INT would 
change the boxed value.
+      record.putValue(LONG_COL, 10_000_000_000L + index * 1_000_000L + i);
+      record.putValue(STR_COL, "s_" + index + "_" + i);
+      // Every 7th row is null so the null-handling test exercises null 
projection.
+      if (i % 7 == 0) {
+        record.addNullValueField(NULLABLE_COL);
+      } else {
+        record.putValue(NULLABLE_COL, i);
+      }
+      records.add(record);
+    }
+    return records;
+  }
+
+  private static List<GenericRow> buildDisjointSortedRecords(int index) {
+    List<GenericRow> records = new ArrayList<>(NUM_RECORDS_PER_SEGMENT);
+    int baseValue = index * 1000;
+    for (int i = 0; i < NUM_RECORDS_PER_SEGMENT; i++) {
+      GenericRow record = new GenericRow();
+      record.putValue(SORTED_COL, baseValue + i);
+      record.putValue(TAIL_COL, i);
+      record.putValue(VAL_COL, baseValue + i);
+      record.putValue(NULLABLE_COL, i);
+      record.putValue(LONG_COL, 10_000_000_000L + baseValue + i);
+      record.putValue(STR_COL, "d_" + index + "_" + i);
+      records.add(record);
+    }
+    return records;
+  }
+
+  private static List<GenericRow> buildLowCardinalityRecords(int index) {
+    List<GenericRow> records = new ArrayList<>(NUM_RECORDS_PER_SEGMENT);
+    for (int i = 0; i < NUM_RECORDS_PER_SEGMENT; i++) {
+      GenericRow record = new GenericRow();
+      // 4 distinct values per segment, non-decreasing so the segment is 
physically sorted on SORTED_COL.
+      record.putValue(SORTED_COL, i / 25);
+      record.putValue(TAIL_COL, NUM_RECORDS_PER_SEGMENT - i);
+      record.putValue(VAL_COL, index * 1_000_000 + i);
+      record.putValue(NULLABLE_COL, i);
+      record.putValue(LONG_COL, 10_000_000_000L + index * 1_000_000L + i);
+      record.putValue(STR_COL, "l_" + index + "_" + i);
+      records.add(record);
+    }
+    return records;
+  }
+
+  private static List<GenericRow> buildUnsortedRecords(int index) {
+    List<GenericRow> records = new ArrayList<>(NUM_RECORDS_PER_SEGMENT);
+    for (int i = 0; i < NUM_RECORDS_PER_SEGMENT; i++) {
+      GenericRow record = new GenericRow();
+      // A non-monotonic permutation of [0, NUM_RECORDS_PER_SEGMENT) (7919 is 
prime and coprime with 100), so the
+      // column is genuinely not physically sorted and SelectionPlanNode falls 
back to a materialized top-K block.
+      record.putValue(SORTED_COL, (i * 7919) % NUM_RECORDS_PER_SEGMENT);
+      record.putValue(TAIL_COL, i);
+      record.putValue(VAL_COL, index * 1_000_000 + i);
+      record.putValue(NULLABLE_COL, i);
+      record.putValue(LONG_COL, 10_000_000_000L + index * 1_000_000L + i);
+      record.putValue(STR_COL, "u_" + index + "_" + i);
+      records.add(record);
+    }
+    return records;
+  }
+
+  private static IndexSegment buildSegment(TableConfig tableConfig, String 
segmentName, List<GenericRow> records,
+      boolean nullHandling)
+      throws Exception {
+    SegmentGeneratorConfig segmentGeneratorConfig = new 
SegmentGeneratorConfig(tableConfig, SCHEMA);
+    segmentGeneratorConfig.setTableName(RAW_TABLE_NAME);
+    segmentGeneratorConfig.setSegmentName(segmentName);
+    segmentGeneratorConfig.setDefaultNullHandlingEnabled(nullHandling);
+    segmentGeneratorConfig.setOutDir(TEMP_DIR.getPath());
+
+    SegmentIndexCreationDriverImpl driver = new 
SegmentIndexCreationDriverImpl();
+    driver.init(segmentGeneratorConfig, new GenericRowRecordReader(records));
+    driver.build();
+
+    return ImmutableSegmentLoader.load(new File(TEMP_DIR, segmentName), 
ReadMode.mmap);
+  }
+
+  @Test
+  public void testAscendingParity() {
+    assertParity(_sortedSegments, "SELECT sortedCol, valCol FROM testTable 
ORDER BY sortedCol, valCol LIMIT 50", false);
+  }
+
+  @Test
+  public void testDescendingParity() {
+    // Reverse order must be allowed for the per-segment forward-scan to 
iterate sortedCol descending; otherwise the
+    // segment falls back to the materialized DESC operator (covered 
separately by testDescIncompatibleFallbackParity).
+    assertParity(_sortedSegments,
+        "SET allowReverseOrder=true; SELECT sortedCol, valCol FROM testTable 
ORDER BY sortedCol DESC, valCol DESC "
+            + "LIMIT 50", false);
+  }
+
+  @Test
+  public void testDescIncompatibleFallbackParity() {
+    // allowReverseOrder=false + DESC -> the streaming child cannot scan 
descending, so SelectionPlanNode emits the
+    // materialized DESC top-K block; the combine still routes to the 
streaming combine and merges single-block cursors.
+    assertParity(_sortedSegments,
+        "SET allowReverseOrder=false; SELECT sortedCol, valCol FROM testTable 
ORDER BY sortedCol DESC, valCol DESC "
+            + "LIMIT 50", false);
+  }
+
+  @Test
+  public void testLimitOffsetParity() {
+    // The server retains limit + offset rows; the broker applies the offset 
later, so both operators keep 40 rows.
+    assertParity(_sortedSegments,
+        "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol 
LIMIT 30 OFFSET 10", false);
+  }
+
+  @Test
+  public void testLowCardinalityMultiColumnParity() {
+    // Low-cardinality primary column => long runs and many sortedCol ties at 
the prune boundary; valCol breaks ties.
+    assertParity(_lowCardSegments,
+        "SELECT sortedCol, tailCol, valCol FROM testTable ORDER BY sortedCol, 
valCol LIMIT 40", false);
+  }
+
+  @Test
+  public void testTwoPhaseSelectNonOrderByParity() {
+    // tailCol is selected but not an order-by key -> the streaming children 
take the two-phase (order-by-then-fetch)
+    // path. valCol is order-by-only, exercising the phase-1 projection of a 
non-selected order-by column.
+    assertParity(_sortedSegments, "SELECT tailCol, sortedCol FROM testTable 
ORDER BY sortedCol, valCol LIMIT 50",
+        false);
+  }
+
+  @Test
+  public void testNonIntProjectionParity() {
+    // Projects a LONG and a STRING column so the multiset comparison would 
catch a stored-type / boxing regression that
+    // an all-INT projection cannot observe.
+    assertParity(_sortedSegments,
+        "SELECT strCol, longCol, sortedCol FROM testTable ORDER BY sortedCol, 
valCol LIMIT 50", false);
+  }
+
+  @Test
+  public void testMixedSortedAndUnsortedSegmentsParity() {
+    assertParity(_mixedSegments, "SELECT sortedCol, valCol FROM testTable 
ORDER BY sortedCol, valCol LIMIT 50", false);
+  }
+
+  @Test
+  public void testNullHandlingEnabledParity() {
+    // Null handling on disables min/max pruning (the combine activates every 
segment); nullableCol carries real nulls.
+    assertParity(_sortedSegments,
+        "SELECT nullableCol, sortedCol, valCol FROM testTable ORDER BY 
sortedCol, valCol LIMIT 50", true);
+  }
+
+  @Test
+  public void testNullHandlingDisabledParity() {
+    // Same segments/query as the enabled case but null handling off: nulls 
read back as the column default, pruning on.
+    assertParity(_sortedSegments,
+        "SELECT nullableCol, sortedCol, valCol FROM testTable ORDER BY 
sortedCol, valCol LIMIT 50", false);
+  }
+
+  @Test
+  public void testStreamingMultiBlockExactCount() {
+    // A precise check on the bounded-flush behavior: 20 rows flushed in 
blocks of 3 yields ceil(20/3) = 7 data blocks.
+    int blockSize = 3;
+    int limit = 20;
+    @Language("sql") String query =
+        "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol 
LIMIT " + limit;
+    Result streaming = run(_sortedSegments, query, true, false, true, 
blockSize);
+    assertEquals(streaming._combineOperator.getClass(), 
StreamingSelectionOrderByCombineOperator.class);
+    assertEquals(streaming._rows.size(), limit);
+    assertEquals(streaming._numBlocks, (limit + blockSize - 1) / blockSize, 
"Unexpected number of streamed blocks");
+    assertSorted(streaming._rows, orderByComparator(query, false));
+    assertMultisetEquals(streaming._rows, run(_sortedSegments, query, false, 
false, false, 0)._rows);
+  }
+
+  @Test
+  public void testPruningSkipsOutOfTopKSegments() {
+    // Disjoint, globally-increasing ranges + small LIMIT: only the lowest 
segment can contribute, the rest are pruned
+    // (never acquired or scanned), so far fewer than all docs are scanned.
+    Result result = run(_disjointSegments, "SELECT sortedCol, valCol FROM 
testTable ORDER BY sortedCol, valCol LIMIT 5",
+        true, false, false, 0);
+    assertTrue(result._combineOperator instanceof 
StreamingSelectionOrderByCombineOperator);
+    assertEquals(result._rows.size(), 5);
+    for (int i = 0; i < 5; i++) {
+      assertEquals((int) result._rows.get(i)[0], i, "Unexpected value at 
position " + i);
+    }
+    int totalDocs = NUM_SEGMENTS * NUM_RECORDS_PER_SEGMENT;
+    assertTrue(result._numDocsScanned < totalDocs, "Pruning should avoid 
scanning every doc, scanned: "
+        + result._numDocsScanned + " of " + totalDocs);
+    assertTrue(result._numDocsScanned <= NUM_RECORDS_PER_SEGMENT,
+        "Only the lowest-range segment should be scanned, but docs scanned 
was: " + result._numDocsScanned);
+  }
+
+  @Test
+  public void testEmptyResultSchemaFallback() {
+    // A filter that matches nothing: every streaming child returns no rows, 
so the combine rebuilds the result schema
+    // from the segment metadata (order-by expressions first) rather than from 
a child block.
+    Result result = run(_sortedSegments,
+        "SELECT sortedCol, valCol FROM testTable WHERE sortedCol < 0 ORDER BY 
sortedCol, valCol LIMIT 10", true, false,
+        false, 0);
+    assertTrue(result._combineOperator instanceof 
StreamingSelectionOrderByCombineOperator);
+    assertTrue(result._rows.isEmpty(), "Expected an empty result, got: " + 
result._rows.size() + " rows");
+    assertEquals(result._schema, new DataSchema(new String[]{SORTED_COL, 
VAL_COL},
+        new ColumnDataType[]{ColumnDataType.INT, ColumnDataType.INT}));
+  }
+
+  @Test
+  public void testHintOffSelectsMinMaxOperator() {
+    // Default behavior is unchanged when the hint is off: the classic MinMax 
operator is still selected.
+    Result result = run(_sortedSegments, "SELECT sortedCol, valCol FROM 
testTable ORDER BY sortedCol, valCol LIMIT 50",
+        false, false, false, 0);
+    assertTrue(result._combineOperator instanceof 
MinMaxValueBasedSelectionOrderByCombineOperator,
+        "Hint off must keep the default MinMax combine operator, got: "
+            + result._combineOperator.getClass().getSimpleName());
+  }
+
+  @Test
+  public void testNonIdentifierOrderByFallsBack() {
+    // Even with the hint on, a non-identifier first order-by expression falls 
back to SelectionOrderByCombineOperator
+    // (the streaming operator and its segment-level counterpart only support 
a leading identifier).
+    Result result = run(_sortedSegments,
+        "SELECT sortedCol, valCol FROM testTable ORDER BY ADD(sortedCol, 1), 
valCol LIMIT 50", true, false, false, 0);
+    assertEquals(result._combineOperator.getClass(), 
SelectionOrderByCombineOperator.class,
+        "Non-identifier first order-by must fall back to 
SelectionOrderByCombineOperator, got: "
+            + result._combineOperator.getClass().getSimpleName());
+  }
+
+  /// Asserts streaming-vs-non-streaming parity for {@code query}. Runs the 
MinMax combine (hint off) as the reference,
+  /// then runs the streaming combine in BOTH single-block mode and bounded 
multi-block streaming mode, asserting each
+  /// selects the streaming operator and produces rows that are sorted by the 
order-by comparator and equal the MinMax
+  /// rows as a multiset. The streaming variant additionally checks the 
bounded-flush invariants.
+  private void assertParity(List<IndexSegment> segments, @Language("sql") 
String query, boolean nullHandling) {
+    Result baseline = run(segments, query, false, nullHandling, false, 0);
+    assertEquals(baseline._combineOperator.getClass(), 
MinMaxValueBasedSelectionOrderByCombineOperator.class,
+        "Baseline must be the MinMax combine operator, got: " + 
baseline._combineOperator.getClass().getSimpleName());
+    Comparator<Object[]> comparator = orderByComparator(query, nullHandling);
+
+    // Classic single-stage path (null streamer): the streaming combine 
flushes the whole merge as one block.
+    Result singleStage = run(segments, query, true, nullHandling, false, 0);
+    assertStreamingParity(singleStage, baseline, comparator, query, false, 0);
+
+    // MSE leaf path (non-null streamer): a small block size forces several 
bounded data blocks before the metadata
+    // block, genuinely exercising the streaming flush path rather than a 
single trimmed block.
+    int blockSize = 3;
+    Result streamed = run(segments, query, true, nullHandling, true, 
blockSize);
+    assertStreamingParity(streamed, baseline, comparator, query, true, 
blockSize);
+  }
+
+  private void assertStreamingParity(Result result, Result baseline, 
Comparator<Object[]> comparator,
+      @Language("sql") String query, boolean streaming, 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);
+      }
+    }
+  }
+
+  /// Runs one combine over {@code segments} and collects its rows, blocks, 
schema and docs-scanned stat.
+  private Result run(List<IndexSegment> segments, @Language("sql") String 
query, boolean hintOn, boolean nullHandling,
+      boolean streaming, int blockSize) {
+    QueryContext queryContext = 
QueryContextConverterUtils.getQueryContext(query);
+    queryContext.setNullHandlingEnabled(nullHandling);
+    if (hintOn) {
+      queryContext.setSortedSelectionMergeEnabled(true);
+      if (blockSize > 0) {
+        queryContext.setSortedSelectionMergeBlockSize(blockSize);
+      }
+    }
+    queryContext.setEndTimeMs(System.currentTimeMillis() + 
Server.DEFAULT_QUERY_EXECUTOR_TIMEOUT_MS);
+
+    List<PlanNode> planNodes = new ArrayList<>(segments.size());
+    for (IndexSegment segment : segments) {
+      SegmentContext segmentContext = new SegmentContext(segment);
+      planNodes.add(streaming ? 
PLAN_MAKER.makeStreamingSegmentPlanNode(segmentContext, queryContext)
+          : PLAN_MAKER.makeSegmentPlanNode(segmentContext, queryContext));
+    }
+    ResultsBlockStreamer streamer = streaming ? block -> {
+    } : null;
+    CombinePlanNode combinePlanNode = new CombinePlanNode(planNodes, 
queryContext, EXECUTOR, streamer);
+
+    Result result = new Result();
+    Operator<?> combineOperator = combinePlanNode.run();
+    result._combineOperator = combineOperator;
+    result._rows = new ArrayList<>();
+    result._blockSizes = new ArrayList<>();
+    if (streaming) {
+      // Drive the streaming combine: collect bounded data blocks until the 
terminal metadata block, which carries the
+      // aggregated execution stats.
+      while (true) {
+        BaseResultsBlock block = (BaseResultsBlock) 
combineOperator.nextBlock();
+        if (block instanceof MetadataResultsBlock) {
+          if (result._schema == null) {
+            result._schema = block.getDataSchema();
+          }
+          result._numDocsScanned = block.getNumDocsScanned();
+          break;
+        }
+        SelectionResultsBlock dataBlock = (SelectionResultsBlock) block;
+        if (result._schema == null) {
+          result._schema = dataBlock.getDataSchema();
+        }
+        List<Object[]> rows = dataBlock.getRows();
+        assertNotNull(rows);
+        result._rows.addAll(rows);
+        result._blockSizes.add(rows.size());
+        result._numBlocks++;
+        assertTrue(result._numBlocks < 1_000_000, "Streaming combine did not 
terminate");
+      }
+    } else {
+      SelectionResultsBlock block = (SelectionResultsBlock) 
combineOperator.nextBlock();
+      result._schema = block.getDataSchema();
+      List<Object[]> rows = block.getRows();
+      assertNotNull(rows);
+      result._rows.addAll(rows);
+      result._blockSizes.add(rows.size());
+      result._numBlocks = 1;
+      result._numDocsScanned = block.getNumDocsScanned();
+    }
+    return result;
+  }
+
+  private static Comparator<Object[]> orderByComparator(@Language("sql") 
String query, boolean nullHandling) {
+    QueryContext queryContext = 
QueryContextConverterUtils.getQueryContext(query);
+    List<OrderByExpressionContext> orderByExpressions = 
queryContext.getOrderByExpressions();
+    assertNotNull(orderByExpressions);
+    return OrderByComparatorFactory.getComparator(orderByExpressions, 
nullHandling);
+  }
+
+  private static void assertSorted(List<Object[]> rows, Comparator<Object[]> 
comparator) {
+    for (int i = 1; i < rows.size(); i++) {
+      assertTrue(comparator.compare(rows.get(i - 1), rows.get(i)) <= 0,
+          "Rows not sorted by the order-by comparator at position " + i);
+    }
+  }
+
+  /// Asserts the two row lists contain the same rows, independent of the 
ordering of fully-equal projected rows.
+  private static void assertMultisetEquals(List<Object[]> actual, 
List<Object[]> expected) {
+    assertEquals(toCanonical(actual), toCanonical(expected), "Row multisets 
differ");
+  }
+
+  /// Canonicalizes rows for multiset comparison. Each cell is encoded with 
its runtime class so a stored-type / boxing
+  /// regression (e.g. a LONG emitted where the reference emits INT) changes 
the encoding and fails the assertion, which
+  /// a plain {@code Arrays.toString} (type-blind) comparison would miss.
+  private static List<String> toCanonical(List<Object[]> rows) {
+    return rows.stream().map(row -> {
+      StringBuilder sb = new StringBuilder("[");
+      for (int i = 0; i < row.length; i++) {
+        if (i > 0) {
+          sb.append(", ");
+        }
+        Object cell = row[i];
+        sb.append(cell == null ? "null" : cell.getClass().getSimpleName() + 
":" + cell);
+      }
+      return sb.append(']').toString();
+    }).sorted().collect(Collectors.toList());
+  }
+
+  @AfterClass
+  public void tearDown()
+      throws IOException {
+    EXECUTOR.shutdownNow();
+    for (List<IndexSegment> segments : List.of(_sortedSegments, 
_disjointSegments, _mixedSegments, _lowCardSegments)) {
+      for (IndexSegment segment : segments) {
+        segment.destroy();
+      }
+    }
+    FileUtils.deleteDirectory(TEMP_DIR);
+  }
+
+  /// Captured output of a single combine run.
+  private static class Result {
+    private Operator<?> _combineOperator;
+    private DataSchema _schema;
+    private List<Object[]> _rows;
+    private List<Integer> _blockSizes;
+    private int _numBlocks;
+    private long _numDocsScanned;
+  }
+}
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
new file mode 100644
index 00000000000..73a071f673a
--- /dev/null
+++ 
b/pinot-core/src/test/java/org/apache/pinot/core/operator/query/StreamingSelectionOrderByOperatorTest.java
@@ -0,0 +1,406 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.core.operator.query;
+
+import java.io.File;
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.List;
+import org.apache.commons.io.FileUtils;
+import org.apache.pinot.common.utils.DataSchema;
+import org.apache.pinot.core.common.Operator;
+import org.apache.pinot.core.operator.blocks.results.SelectionResultsBlock;
+import org.apache.pinot.core.plan.SelectionPlanNode;
+import org.apache.pinot.core.query.request.context.QueryContext;
+import 
org.apache.pinot.core.query.request.context.utils.QueryContextConverterUtils;
+import 
org.apache.pinot.segment.local.indexsegment.immutable.ImmutableSegmentLoader;
+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.IndexSegment;
+import org.apache.pinot.segment.spi.SegmentContext;
+import org.apache.pinot.segment.spi.creator.SegmentGeneratorConfig;
+import org.apache.pinot.spi.config.table.TableConfig;
+import org.apache.pinot.spi.config.table.TableType;
+import org.apache.pinot.spi.data.FieldSpec;
+import org.apache.pinot.spi.data.Schema;
+import org.apache.pinot.spi.data.readers.GenericRow;
+import org.apache.pinot.spi.utils.ReadMode;
+import org.apache.pinot.spi.utils.builder.TableConfigBuilder;
+import org.intellij.lang.annotations.Language;
+import org.testng.annotations.AfterClass;
+import org.testng.annotations.BeforeClass;
+import org.testng.annotations.Test;
+
+import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertFalse;
+import static org.testng.Assert.assertNotNull;
+import static org.testng.Assert.assertTrue;
+
+
+/// Segment-level tests for {@link StreamingSelectionOrderByOperator}.
+///
+/// <p>The operator emits the same globally-sorted rows as the existing 
materialized selection ORDER BY operators, only
+/// spread across many lazily-produced blocks. Each test therefore asserts 
<b>stream-vs-materialized parity</b>: it
+/// drives the streaming operator through {@link SelectionPlanNode} with 
{@code sortedSelectionMergeEnabled=true},
+/// concatenates
+/// every {@link Operator#nextBlock()} output until {@code null}, and asserts 
the concatenation equals the single block
+/// the materialized operator ({@link 
SelectionPartiallyOrderedByLinearOperator} / {@link SelectionOrderByOperator})
+/// produces for the identical query with the hint off.
+///
+/// <p>To keep element-wise comparison deterministic (priority-queue draining 
is not stable for rows that tie on every
+/// order-by column) each fixture makes the order-by column tuple unique per 
row: the {@code _segment} fixture has a
+/// unique sorted column, while the {@code _dupSegment} / {@code 
_largeSegment} fixtures repeat the sorted column but
+/// pair it with a unique {@code TAIL_COL} order-by tail. The 
all-order-by-columns-tie case (where stream and
+/// materialized output may legitimately differ in row order) is intentionally 
out of scope here and is covered by the
+/// combine-level test via multiset comparison.
+public class StreamingSelectionOrderByOperatorTest {
+  private static final File TEMP_DIR = new File(FileUtils.getTempDirectory(), 
"StreamingSelectionOrderByOperatorTest");
+  private static final String RAW_TABLE_NAME = "testTable";
+
+  private static final String SORTED_COL = "sortedCol";
+  private static final String TAIL_COL = "tailCol";
+  private static final String VAL_COL = "valCol";
+  private static final String NULLABLE_COL = "nullableCol";
+
+  private static final TableConfig TABLE_CONFIG =
+      new 
TableConfigBuilder(TableType.OFFLINE).setTableName(RAW_TABLE_NAME).setSortedColumn(SORTED_COL).build();
+  private static final Schema SCHEMA = new Schema.SchemaBuilder()
+      .addSingleValueDimension(SORTED_COL, FieldSpec.DataType.INT)
+      .addSingleValueDimension(TAIL_COL, FieldSpec.DataType.INT)
+      .addSingleValueDimension(VAL_COL, FieldSpec.DataType.INT)
+      .addSingleValueDimension(NULLABLE_COL, FieldSpec.DataType.INT)
+      .build();
+
+  /// Unique sorted column, no nulls. Exercises the no-tail emission mode.
+  private static final int NUM_RECORDS = 30;
+  /// Repeated sorted column (RUN_SIZE rows per value) paired with a unique 
tail. Exercises the run / heap path, and
+  /// carries nulls in NULLABLE_COL for the null-handling cases.
+  private static final int NUM_DISTINCT_SORTED = 10;
+  private static final int RUN_SIZE = 4;
+  private static final int NUM_DUP_RECORDS = NUM_DISTINCT_SORTED * RUN_SIZE;
+  /// A single primary-value run larger than one project block 
(DocIdSetPlanNode.MAX_DOC_PER_CALL == 10_000), so the
+  /// forward scan crosses a block boundary within one run.
+  private static final int NUM_LARGE_RECORDS = 12_000;
+  /// Two runs where the first is larger than one project block: with a limit 
greater than the first run's size the
+  /// operator both crosses a project-block boundary mid-run AND carries 
_pendingRow across a run boundary to emit a
+  /// second block.
+  private static final int MULTI_RUN0_SIZE = 10_500;
+  private static final int MULTI_RUN1_SIZE = 1_000;
+  private static final int NUM_MULTI_RUN_RECORDS = MULTI_RUN0_SIZE + 
MULTI_RUN1_SIZE;
+
+  private IndexSegment _segment;
+  private IndexSegment _dupSegment;
+  private IndexSegment _largeSegment;
+  private IndexSegment _multiRunSegment;
+  /// Carries nulls in an order-by (tail) column so the streaming operator's 
own null path (phase-1 materialization and
+  /// the null-aware comparator) is exercised, not just the shared two-phase 
fetch.
+  private IndexSegment _nullTailSegment;
+
+  @BeforeClass
+  public void setUp()
+      throws Exception {
+    FileUtils.deleteDirectory(TEMP_DIR);
+    _segment = buildSegment("uniqueSorted", buildUniqueSortedRecords(), false);
+    _dupSegment = buildSegment("dupSorted", buildDupSortedRecords(), true);
+    _largeSegment = buildSegment("largeRun", buildLargeRunRecords(), false);
+    _multiRunSegment = buildSegment("multiRun", buildMultiRunRecords(), false);
+    _nullTailSegment = buildSegment("nullTail", buildNullTailRecords(), true);
+  }
+
+  private static List<GenericRow> buildUniqueSortedRecords() {
+    List<GenericRow> records = new ArrayList<>(NUM_RECORDS);
+    for (int i = 0; i < NUM_RECORDS; i++) {
+      GenericRow record = new GenericRow();
+      record.putValue(SORTED_COL, i);
+      record.putValue(TAIL_COL, NUM_RECORDS - i);
+      record.putValue(VAL_COL, i * 3);
+      record.putValue(NULLABLE_COL, i);
+      records.add(record);
+    }
+    return records;
+  }
+
+  private static List<GenericRow> buildDupSortedRecords() {
+    List<GenericRow> records = new ArrayList<>(NUM_DUP_RECORDS);
+    for (int i = 0; i < NUM_DUP_RECORDS; i++) {
+      GenericRow record = new GenericRow();
+      record.putValue(SORTED_COL, i / RUN_SIZE);
+      // Tail resets to a descending sequence within each run so the column is 
NOT globally sorted (the run/heap path is
+      // only taken when the tail is unsorted); the (sortedCol, tailCol) tuple 
is still unique per row.
+      record.putValue(TAIL_COL, RUN_SIZE - 1 - (i % RUN_SIZE));
+      record.putValue(VAL_COL, i * 2);
+      // Every third row carries a null so the two-phase fetch is exercised 
with and without null handling.
+      if (i % 3 == 0) {
+        record.addNullValueField(NULLABLE_COL);
+      } else {
+        record.putValue(NULLABLE_COL, i);
+      }
+      records.add(record);
+    }
+    return records;
+  }
+
+  private static List<GenericRow> buildLargeRunRecords() {
+    List<GenericRow> records = new ArrayList<>(NUM_LARGE_RECORDS);
+    for (int i = 0; i < NUM_LARGE_RECORDS; i++) {
+      GenericRow record = new GenericRow();
+      // All rows share one sorted value, so they form a single run spanning 
multiple project blocks.
+      record.putValue(SORTED_COL, 0);
+      // A non-monotonic permutation of [0, NUM_LARGE_RECORDS) (7919 is prime 
and coprime with NUM_LARGE_RECORDS, so the
+      // mapping is a bijection): unique tail values that are not physically 
sorted, forcing the run/heap path.
+      record.putValue(TAIL_COL, (i * 7919) % NUM_LARGE_RECORDS);
+      record.putValue(VAL_COL, i);
+      record.putValue(NULLABLE_COL, i);
+      records.add(record);
+    }
+    return records;
+  }
+
+  private static List<GenericRow> buildMultiRunRecords() {
+    List<GenericRow> records = new ArrayList<>(NUM_MULTI_RUN_RECORDS);
+    appendRun(records, 0, MULTI_RUN0_SIZE);
+    appendRun(records, 1, MULTI_RUN1_SIZE);
+    return records;
+  }
+
+  /// Appends one run of {@code runSize} rows all sharing {@code sortedValue}, 
with a tail that descends within the run
+  /// (so the tail column is not globally sorted) and is unique per 
(sortedCol, tailCol) tuple.
+  private static void appendRun(List<GenericRow> records, int sortedValue, int 
runSize) {
+    for (int j = 0; j < runSize; j++) {
+      GenericRow record = new GenericRow();
+      record.putValue(SORTED_COL, sortedValue);
+      record.putValue(TAIL_COL, runSize - 1 - j);
+      record.putValue(VAL_COL, sortedValue * 1_000_000 + j);
+      record.putValue(NULLABLE_COL, j);
+      records.add(record);
+    }
+  }
+
+  private static List<GenericRow> buildNullTailRecords() {
+    List<GenericRow> records = new ArrayList<>();
+    // Three runs, each with exactly one null tail and two distinct non-null 
tails, so ordering stays deterministic
+    // (no all-order-by-columns tie) while a null flows through the order-by 
tail column.
+    for (int g = 0; g < 3; g++) {
+      for (int j = 0; j < 3; j++) {
+        GenericRow record = new GenericRow();
+        record.putValue(SORTED_COL, g);
+        if (j == 0) {
+          record.addNullValueField(TAIL_COL);
+        } else {
+          record.putValue(TAIL_COL, g * 10 + j);
+        }
+        record.putValue(VAL_COL, g * 100 + j);
+        record.putValue(NULLABLE_COL, g * 100 + j);
+        records.add(record);
+      }
+    }
+    return records;
+  }
+
+  private static IndexSegment buildSegment(String segmentName, 
List<GenericRow> records, boolean nullHandling)
+      throws Exception {
+    SegmentGeneratorConfig segmentGeneratorConfig = new 
SegmentGeneratorConfig(TABLE_CONFIG, SCHEMA);
+    segmentGeneratorConfig.setTableName(RAW_TABLE_NAME);
+    segmentGeneratorConfig.setSegmentName(segmentName);
+    segmentGeneratorConfig.setDefaultNullHandlingEnabled(nullHandling);
+    segmentGeneratorConfig.setOutDir(TEMP_DIR.getPath());
+
+    SegmentIndexCreationDriverImpl driver = new 
SegmentIndexCreationDriverImpl();
+    driver.init(segmentGeneratorConfig, new GenericRowRecordReader(records));
+    driver.build();
+
+    return ImmutableSegmentLoader.load(new File(TEMP_DIR, segmentName), 
ReadMode.mmap);
+  }
+
+  @Test
+  public void testSingleSortedColumnAscending() {
+    assertParity(_segment, "SELECT sortedCol FROM testTable ORDER BY 
sortedCol", false, 1);
+  }
+
+  @Test
+  public void testSingleSortedColumnDescending() {
+    // Reverse order must be allowed for the forward-scan project to iterate 
the sorted column descending; otherwise
+    // SelectionPlanNode falls back to the materialized DESC operator and the 
streaming operator is never built.
+    assertParity(_segment, "SET allowReverseOrder=true; SELECT sortedCol FROM 
testTable ORDER BY sortedCol DESC", false,
+        1);
+  }
+
+  @Test
+  public void testSortedPrefixWithUnsortedTail() {
+    // Repeated sorted value + unique tail exercises nextRun(): per-run top-K 
heap and the one-row lookahead. With the
+    // default LIMIT 10 spanning ~3 runs of 4, the streaming operator emits 
more than one block.
+    assertParity(_dupSegment, "SELECT sortedCol, tailCol FROM testTable ORDER 
BY sortedCol, tailCol", false, 2);
+  }
+
+  @Test
+  public void testTwoPhaseSingleSortedColumn() {
+    // Output has a non-order-by column (valCol) -> two-phase fetch. Unique 
sorted column keeps ordering deterministic.
+    assertParity(_segment, "SELECT valCol, sortedCol FROM testTable ORDER BY 
sortedCol", false, 1);
+  }
+
+  @Test
+  public void testTwoPhaseWithUnsortedTail() {
+    // Two-phase fetch combined with the run/heap path (sorted prefix + 
unsorted tail).
+    assertParity(_dupSegment, "SELECT valCol FROM testTable ORDER BY 
sortedCol, tailCol", false, 1);
+  }
+
+  @Test
+  public void testLimitOffsetNoTail() {
+    // Server retains limit + offset rows; the broker applies the offset later.
+    assertParity(_segment, "SELECT sortedCol FROM testTable ORDER BY sortedCol 
LIMIT 5 OFFSET 3", false, 1);
+  }
+
+  @Test
+  public void testLimitOffsetWithTail() {
+    assertParity(_dupSegment, "SELECT sortedCol, tailCol FROM testTable ORDER 
BY sortedCol, tailCol LIMIT 7 OFFSET 5",
+        false, 1);
+  }
+
+  @Test
+  public void testRunSpanningMultipleProjectBlocks() {
+    // One run of 12_000 rows (> one 10k project block); the heap caps at 
limit + offset while the scan crosses the
+    // block boundary. A single primary value means a single run, hence a 
single emitted block.
+    StreamingSelectionOrderByOperator operator = assertParity(_largeSegment,
+        "SELECT sortedCol, tailCol FROM testTable ORDER BY sortedCol, tailCol 
LIMIT 25", false, 1);
+    // The whole run is scanned to find its top-K, which only happens if the 
forward scan pulled every project block
+    // (proving the scan genuinely crossed the 10k boundary rather than 
stopping at the first block).
+    assertEquals(operator.getExecutionStatistics().getNumDocsScanned(), 
NUM_LARGE_RECORDS);
+  }
+
+  @Test
+  public void testRunBoundaryAcrossProjectBlocks() {
+    // First run (10_500 rows) is larger than one project block, and the limit 
(10_600) exceeds it, so the operator
+    // crosses a project-block boundary mid-run AND carries _pendingRow across 
the run boundary to emit a second block.
+    StreamingSelectionOrderByOperator operator = assertParity(_multiRunSegment,
+        "SELECT sortedCol, tailCol FROM testTable ORDER BY sortedCol, tailCol 
LIMIT 10600", false, 2);
+    assertEquals(operator.getExecutionStatistics().getNumDocsScanned(), 
NUM_MULTI_RUN_RECORDS);
+  }
+
+  @Test
+  public void testNullHandlingEnabled() {
+    assertParity(_dupSegment, "SELECT nullableCol, sortedCol, tailCol FROM 
testTable ORDER BY sortedCol, tailCol", true,
+        1);
+  }
+
+  @Test
+  public void testNullHandlingDisabled() {
+    assertParity(_dupSegment, "SELECT nullableCol, sortedCol, tailCol FROM 
testTable ORDER BY sortedCol, tailCol",
+        false, 1);
+  }
+
+  @Test
+  public void testNullInOrderByColumnWithNullHandling() {
+    // Drives a null through the order-by tail column (not just the carried 
non-order-by column), exercising the
+    // streaming operator's null-aware comparator and phase-1 null 
materialization.
+    String query = "SELECT tailCol, sortedCol, valCol FROM testTable ORDER BY 
sortedCol, tailCol";
+    assertParity(_nullTailSegment, query, true, 1);
+    // Absolute anchor: with null handling on, the order-by tail column (index 
1 after extractExpressions) actually
+    // carries nulls through to the output. Guards against both operators 
substituting a default value identically.
+    List<Object[]> rows = collectStreamingRows(_nullTailSegment, query, true);
+    assertTrue(rows.stream().anyMatch(row -> row[1] == null), "Expected a null 
in the order-by tail column");
+  }
+
+  @Test
+  public void testNullInOrderByColumnWithoutNullHandling() {
+    String query = "SELECT tailCol, sortedCol, valCol FROM testTable ORDER BY 
sortedCol, tailCol";
+    assertParity(_nullTailSegment, query, false, 1);
+    // With null handling off the null reads back as the column's default 
value, so no output cell is null.
+    List<Object[]> rows = collectStreamingRows(_nullTailSegment, query, false);
+    assertTrue(rows.stream().noneMatch(row -> row[1] == null),
+        "Expected no nulls in the output when null handling is disabled");
+  }
+
+  /// Runs {@code query} twice over {@code 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.
+  ///
+  /// @param expectedMinBlocks the minimum number of non-null blocks the 
streaming operator must emit. Tail-mode
+  ///     cases that span multiple runs pass {@code >= 2} to prove the output 
is genuinely streamed; cases whose
+  ///     result fits in a single trimmed block pass {@code 1}.
+  private StreamingSelectionOrderByOperator assertParity(IndexSegment segment, 
@Language("sql") String query,
+      boolean nullHandling, int expectedMinBlocks) {
+    // Streaming path.
+    QueryContext streamingContext = 
QueryContextConverterUtils.getQueryContext(query);
+    streamingContext.setNullHandlingEnabled(nullHandling);
+    streamingContext.setSortedSelectionMergeEnabled(true);
+    Operator<SelectionResultsBlock> streamingOperator =
+        new SelectionPlanNode(new SegmentContext(segment), 
streamingContext).run();
+    assertTrue(streamingOperator instanceof StreamingSelectionOrderByOperator,
+        "Expected the streaming operator to be built, got: " + 
streamingOperator.getClass().getSimpleName());
+
+    List<Object[]> streamingRows = new ArrayList<>();
+    DataSchema streamingSchema = null;
+    int numBlocks = 0;
+    SelectionResultsBlock block;
+    while ((block = streamingOperator.nextBlock()) != null) {
+      numBlocks++;
+      if (streamingSchema == null) {
+        streamingSchema = block.getDataSchema();
+      }
+      streamingRows.addAll(block.getRows());
+    }
+    assertTrue(numBlocks >= expectedMinBlocks,
+        "Expected at least " + expectedMinBlocks + " streaming block(s), got: 
" + numBlocks);
+
+    // Materialized path (hint off).
+    QueryContext materializedContext = 
QueryContextConverterUtils.getQueryContext(query);
+    materializedContext.setNullHandlingEnabled(nullHandling);
+    Operator<SelectionResultsBlock> materializedOperator =
+        new SelectionPlanNode(new SegmentContext(segment), 
materializedContext).run();
+    assertFalse(materializedOperator instanceof 
StreamingSelectionOrderByOperator,
+        "Materialized baseline must not be the streaming operator");
+    SelectionResultsBlock materializedBlock = materializedOperator.nextBlock();
+    assertNotNull(materializedBlock);
+    List<Object[]> expectedRows = materializedBlock.getRows();
+
+    assertEquals(streamingSchema, materializedBlock.getDataSchema(), "Schema 
mismatch for query: " + query);
+    assertEquals(streamingRows.size(), expectedRows.size(), "Row count 
mismatch for query: " + query);
+    for (int i = 0; i < expectedRows.size(); i++) {
+      assertEquals(streamingRows.get(i), expectedRows.get(i), "Row " + i + " 
mismatch for query: " + query);
+    }
+    return (StreamingSelectionOrderByOperator) streamingOperator;
+  }
+
+  /// Drains the streaming operator for {@code query} and returns all rows 
concatenated in emission order.
+  private List<Object[]> collectStreamingRows(IndexSegment segment, 
@Language("sql") String query,
+      boolean nullHandling) {
+    QueryContext queryContext = 
QueryContextConverterUtils.getQueryContext(query);
+    queryContext.setNullHandlingEnabled(nullHandling);
+    queryContext.setSortedSelectionMergeEnabled(true);
+    Operator<SelectionResultsBlock> operator = new SelectionPlanNode(new 
SegmentContext(segment), queryContext).run();
+    List<Object[]> rows = new ArrayList<>();
+    SelectionResultsBlock block;
+    while ((block = operator.nextBlock()) != null) {
+      rows.addAll(block.getRows());
+    }
+    return rows;
+  }
+
+  @AfterClass
+  public void tearDown()
+      throws IOException {
+    for (IndexSegment segment : new IndexSegment[]{_segment, _dupSegment, 
_largeSegment, _multiRunSegment,
+        _nullTailSegment}) {
+      if (segment != null) {
+        segment.destroy();
+      }
+    }
+    FileUtils.deleteDirectory(TEMP_DIR);
+  }
+}
diff --git 
a/pinot-spi/src/main/java/org/apache/pinot/spi/utils/CommonConstants.java 
b/pinot-spi/src/main/java/org/apache/pinot/spi/utils/CommonConstants.java
index 7c5db4ea1f2..0af0e797710 100644
--- a/pinot-spi/src/main/java/org/apache/pinot/spi/utils/CommonConstants.java
+++ b/pinot-spi/src/main/java/org/apache/pinot/spi/utils/CommonConstants.java
@@ -630,6 +630,11 @@ public class CommonConstants {
     public static final String 
CONFIG_OF_MSE_STREAMING_DISTINCT_FLUSH_THRESHOLD =
         "pinot.broker.mse.streaming.distinct.flush.threshold";
     public static final int DEFAULT_MSE_STREAMING_DISTINCT_FLUSH_THRESHOLD = 
-1;
+
+    /// Default output block size (rows) for the streaming selection ORDER BY 
combine
+    /// ({@link Request.QueryOptionKey#SORTED_SELECTION_MERGE_BLOCK_SIZE}).
+    public static final int DEFAULT_SORTED_SELECTION_MERGE_BLOCK_SIZE = 10_000;
+
     // Whether to infer partition hint by default or not.
     // This value can always be overridden by INFER_PARTITION_HINT query option
     public static final String CONFIG_OF_INFER_PARTITION_HINT = 
"pinot.broker.multistage.infer.partition.hint";
@@ -925,6 +930,11 @@ public class CommonConstants {
         /// there is no such stage and the client would observe duplicate rows 
across flush windows.
         public static final String STREAMING_DISTINCT_FLUSH_THRESHOLD = 
"streamingDistinctFlushThreshold";
 
+        /// Opt-in: use the streaming k-way-merge selection ORDER BY combine 
over sorted segments.
+        public static final String SORTED_SELECTION_MERGE_ENABLED = 
"sortedSelectionMergeEnabled";
+        /// Output block size (rows) for the streaming selection ORDER BY 
combine.
+        public static final String SORTED_SELECTION_MERGE_BLOCK_SIZE = 
"sortedSelectionMergeBlockSize";
+
         public static final String NUM_REPLICA_GROUPS_TO_QUERY = 
"numReplicaGroupsToQuery";
         public static final String ORDERED_PREFERRED_POOLS = 
"orderedPreferredPools";
         public static final String USE_FIXED_REPLICA = "useFixedReplica";


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

Reply via email to