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]
