rohityadav1993 commented on code in PR #19120:
URL: https://github.com/apache/pinot/pull/19120#discussion_r4125618903


##########
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() {

Review Comment:
   Done in e3a908aca5. Added a comment on pullBlock(): it drops the same unit 
as the merger, where a block is a whole segment. Since a reload swaps the whole 
segment, a mismatch normally starts at the first block and both drop 
everything; they only differ on a mid-segment block, where keeping a clean 
prefix is better than a prefix and a suffix with a hole.
   
   testSchemaMismatchDropsTheRestOfTheSegmentAndIsReportedOnce makes two 
segments diverge mid-stream and checks the kept rows, the untouched segments, 
sorted output, and a single deduplicated error.



##########
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() {

Review Comment:
   Both covered in e3a908aca5. Mismatch: see [14].
   
   Lifecycle: instead of a mock segment, each real leaf is wrapped in a 
counting AcquireReleaseColumnsSegmentOperator (as the prefetch path plans it), 
so the combine's real calls are exercised. Tests check every acquire is 
released exactly once on full drain, on limit with cursors open, on stop() 
(including repeated), and on a child Exception or Error. Each was watched 
failing with its release removed.



##########
pinot-spi/src/main/java/org/apache/pinot/spi/utils/CommonConstants.java:
##########
@@ -546,6 +546,9 @@ public static class Broker {
     public static final String 
CONFIG_OF_MSE_STREAMING_GROUP_BY_FLUSH_THRESHOLD =
         "pinot.broker.mse.streaming.group.by.flush.threshold";
     public static final int DEFAULT_MSE_STREAMING_GROUP_BY_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;

Review Comment:
   1. Reconciled in the description; the default was always 10000.
   
   2. Deferred. The bench harness 
(https://github.com/rohityadav1993/pinot/pull/1) pins 10000 today, and the 
trade-off is best measured with PR2's receiver (#19121) downstream.
   
   3. Done in e3a908aca5 as server config: 
pinot.server.query.executor.sorted.selection.merge.block.size, plus 
...auto.min.sorted.ratio for AUTO, validated at init and applied when the query 
does not set the option (like numGroupsLimit). Server-side, so it reaches the 
MSE leaf and gRPC streaming without broker changes; the constant moved to 
Server. Kept separate from other block sizes since it sizes the blocks streamed 
downstream, not the segment scan.



##########
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.

Review Comment:
   Done in e3a908aca5, one sweep across the PR. Comments only.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to