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 a8065e5459d30ee62f9d7b2ad091f7aee580a622 Author: rohity <[email protected]> AuthorDate: Sun Sep 27 14:07:54 2026 +0000 Address second-round review on the streaming selection ORDER BY combine Test segment release and schema-mismatch handling. Nothing exercised either path before: the existing tests plan the combine from plain segment plan nodes, so the acquire/release calls were no-ops and a leaked segment would have passed, and every segment shared one schema, so the mismatch branch never ran. The new tests wrap each real streaming leaf in a counting AcquireReleaseColumnsSegmentOperator, the way the prefetch path plans it, and assert every acquired segment is released exactly once when the merge drains, when the limit stops it with cursors still open, on stop() mid-merge (including a repeated stop()), and when a child throws an exception or an Error. A mismatched schema on one block must drop the rest of that segment, leave every other segment whole, and report one deduplicated MERGE_RESPONSE error. Each was checked against an operator with the corresponding release or drop removed. Document why the cursor drops the rest of the segment on a mismatch rather than one block: it is the same unit the merger drops, since there a block is a segment's whole result. Add server config defaults for the merge options. sortedSelectionMergeBlockSize and sortedSelectionMergeAutoMinSortedRatio could only be changed per query. Add pinot.server.query.executor.sorted.selection.merge.block.size and pinot.server.query.executor.sorted.selection.merge.auto.min.sorted.ratio, read and validated in InstancePlanMakerImplV2#init and applied whenever the query does not set the option, as numGroupsLimit already is. The block size default moves from Broker to Server, since only the server reads it. applyQueryOptions now always writes both values, so the AUTO threshold tests set the ratio through the query option instead of the query context setter. Also bound CombineSlowOperatorsTest's deadline test with a timeout: a regression there drives a SlowOperator that sleeps for an hour, which would hang the build instead of failing it. Expose the combine's settings in the explain plan. It now reports block size, frontier pruning, tie deferral and the segment / sorted-segment counts as explain attributes, so an MSE explain with explainAskingServers shows which path ran and whether pruning was active. Tests pin that numSegmentsMatched excludes never-activated segments. Use markdown syntax in the /// javadoc added by this change: [Foo], backticks, **bold** and blank-line paragraphs instead of {@link}, {@code}, <b> and <p>, matching the rest of the module. --- .../StreamingSelectionOrderByCombineOperator.java | 129 ++++--- .../query/StreamingSelectionOrderByOperator.java | 48 +-- .../core/plan/maker/InstancePlanMakerImplV2.java | 32 +- .../core/query/request/context/QueryContext.java | 3 +- .../operator/combine/CombineSlowOperatorsTest.java | 9 +- ...reamingSelectionOrderByCombineOperatorTest.java | 381 +++++++++++++++++++-- .../StreamingSelectionOrderByOperatorTest.java | 68 ++-- .../plan/maker/InstancePlanMakerImplV2Test.java | 67 +++- .../apache/pinot/spi/utils/CommonConstants.java | 19 +- 9 files changed, 601 insertions(+), 155 deletions(-) diff --git a/pinot-core/src/main/java/org/apache/pinot/core/operator/combine/StreamingSelectionOrderByCombineOperator.java b/pinot-core/src/main/java/org/apache/pinot/core/operator/combine/StreamingSelectionOrderByCombineOperator.java index 75560d3715b..5597a06bcfc 100644 --- a/pinot-core/src/main/java/org/apache/pinot/core/operator/combine/StreamingSelectionOrderByCombineOperator.java +++ b/pinot-core/src/main/java/org/apache/pinot/core/operator/combine/StreamingSelectionOrderByCombineOperator.java @@ -32,6 +32,7 @@ import org.apache.pinot.common.request.context.OrderByExpressionContext; import org.apache.pinot.common.utils.DataSchema; import org.apache.pinot.core.common.Operator; import org.apache.pinot.core.operator.AcquireReleaseColumnsSegmentOperator; +import org.apache.pinot.core.operator.ExplainAttributeBuilder; 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; @@ -49,17 +50,17 @@ 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 +/// It performs an incremental k-way heap merge across the per-segment operators in `_operators`, returning +/// globally sorted rows in bounded blocks. Each segment is exposed through a [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. +/// [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 +/// [SelectionResultsBlock]-producing operator such as `SelectionOrderByOperator`); the cursor reads that /// one block and iterates its rows. /// -/// A {@link PriorityQueue} of {@link SegmentCursor} ordered by the {@link OrderByComparatorFactory} comparator on +/// A [PriorityQueue] of [SegmentCursor] ordered by the [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 @@ -67,39 +68,43 @@ import org.slf4j.LoggerFactory; /// 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(SegmentCursor)} for the +/// (ASC) / max value (DESC) reusing the `MinMaxValueContext` idea from +/// [MinMaxValueBasedSelectionOrderByCombineOperator]. A cursor is only activated (its segment acquired and first +/// block read) when the merge frontier reaches its min/max, so once `limit + offset` rows are emitted the +/// remaining segments are never acquired or read. See [#activateEligibleCursors(SegmentCursor)] for the /// correctness argument. /// Pruning is disabled when null handling is enabled (an unsorted segment's first order-by column may then contain /// nulls whose ordering position the raw min/max cannot capture), in which case every segment is activated. +/// Its effect shows in the stats: a never-activated segment scans no docs, so `numSegmentsMatched` excludes it while +/// `numSegmentsProcessed` counts every segment. Their difference also counts activated segments that matched no rows. +/// TODO: report the activated count separately; that needs a new `DataTable.MetadataKey`. +/// See https://github.com/apache/pinot/pull/19120#discussion_r3871714063 /// /// **Segment acquire/release lifecycle.** A cursor acquires its -/// {@link AcquireReleaseColumnsSegmentOperator} on activation and releases it only when its child operator is fully +/// [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()} +/// [StreamingSelectionOrderByOperator] retains a buffer-backed `ValueBlock` across `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()}. +/// The rows handed out by the child operators are already deep-copied to heap `Object[]` (via +/// `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 [#releaseAllCursors()]. /// /// **Streaming paths only.** This operator is installed only where a -/// {@link org.apache.pinot.core.query.executor.ResultsBlockStreamer} is present (the -/// MSE leaf, driven by {@link org.apache.pinot.core.operator.streaming.StreamingInstanceResponseOperator}, and the -/// gRPC streaming endpoint). The merge emits many bounded {@link SelectionResultsBlock}s from successive -/// {@link #getNextBlock()} calls followed by a final {@link MetadataResultsBlock}. On the blocking single-stage -/// path the caller expects one complete block from a single {@link #getNextBlock()} call, so nothing here could -/// stream: the merge would accumulate the whole {@code limit + offset} result while holding every activated cursor's -/// materialized block alive. That path keeps {@link MinMaxValueBasedSelectionOrderByCombineOperator}. +/// [org.apache.pinot.core.query.executor.ResultsBlockStreamer] is present (the +/// MSE leaf, driven by [org.apache.pinot.core.operator.streaming.StreamingInstanceResponseOperator], and the +/// gRPC streaming endpoint). The merge emits many bounded [SelectionResultsBlock]s from successive +/// [#getNextBlock()] calls followed by a final [MetadataResultsBlock]. On the blocking single-stage +/// path the caller expects one complete block from a single [#getNextBlock()] call, so nothing here could +/// stream: the merge would accumulate the whole `limit + offset` result while holding every activated cursor's +/// materialized block alive. That path keeps [MinMaxValueBasedSelectionOrderByCombineOperator]. /// -/// **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. The base entry points that it replaces - {@link #processSegments()} and -/// {@link #isQuerySatisfied(SelectionResultsBlock, Object)} - are overridden to fail loud. The base -/// {@code Phaser} (which exists only to fence worker threads against segment release) is intentionally bypassed because +/// **Threading.** This operator overrides [#start()]/[#stop()] to no-ops (other than releasing +/// segments) and runs the merge single-threaded and lazily in [#getNextBlock()] on the consumer thread; it does +/// not use the base worker-queue model. The base entry points that it replaces - [#processSegments()] and +/// [#isQuerySatisfied(SelectionResultsBlock, Object)] - are overridden to fail loud. The base +/// `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. @@ -111,10 +116,12 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi private final boolean _asc; private final boolean _pruningEnabled; /// Whether a cursor whose bound *ties* the merge frontier may be deferred, not only one sorting strictly - /// beyond it. Only correct under a single order-by expression: see {@link #sortsBeyond}. + /// beyond it. Only correct under a single order-by expression: see [#sortsBeyond]. private final boolean _deferTiedCursors; private final int _numRowsToKeep; private final int _blockSize; + /// Segments whose metadata marks the leading order-by column sorted; -1 when that expression is not a column + private final int _numSortedSegments; private final Comparator<Object[]> _comparator; private final SegmentCursor[] _sortedCursors; private final PriorityQueue<SegmentCursor> _priorityQueue; @@ -168,6 +175,7 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi // Reading DataSourceMetadata does not touch column buffers, so no segment acquire is needed here (mirrors // MinMaxValueBasedSelectionOrderByCombineOperator). _sortedCursors = new SegmentCursor[_numOperators]; + int numSortedSegments = firstOrderByColumn != null ? 0 : -1; for (int i = 0; i < _numOperators; i++) { Operator<BaseResultsBlock> operator = _operators.get(i); if (firstOrderByColumn == null) { @@ -178,7 +186,11 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi operator.getIndexSegment().getDataSource(firstOrderByColumn, queryContext.getSchema()) .getDataSourceMetadata(); _sortedCursors[i] = new SegmentCursor(operator, metadata.getMinValue(), metadata.getMaxValue()); + if (metadata.isSorted()) { + numSortedSegments++; + } } + _numSortedSegments = numSortedSegments; if (firstOrderByColumn != null) { sortCursorsByMinMax(); } @@ -190,7 +202,7 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi /// 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}). + /// always be processed (mirrors [MinMaxValueBasedSelectionOrderByCombineOperator]). private void sortCursorsByMinMax() { if (_asc) { Arrays.sort(_sortedCursors, (o1, o2) -> { @@ -220,7 +232,23 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi 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 + /// Settings fixed at construction. Explain never drives the merge, so how many segments activate is not shown here. + /// Query-wide values are idempotent so the broker merges the node across workers only when they agree; the segment + /// counts are summed. `numSortedSegments` is metadata sortedness, not whether a child streams: a DESC scan without + /// `allowReverseOrder`, or a leading column holding nulls, still falls back to a materialized child. + @Override + protected void explainAttributes(ExplainAttributeBuilder attributeBuilder) { + super.explainAttributes(attributeBuilder); + attributeBuilder.putLongIdempotent("blockSize", _blockSize); + attributeBuilder.putBool("frontierPruning", _pruningEnabled); + attributeBuilder.putBool("deferTiedCursors", _deferTiedCursors); + attributeBuilder.putLong("numSegments", _numOperators); + if (_numSortedSegments >= 0) { + attributeBuilder.putLong("numSortedSegments", _numSortedSegments); + } + } + + /// Override to a no-op: the merge is single-threaded and lazy in [#getNextBlock()], so we do not spin up the /// base worker threads / blocking-queue model. @Override public void start() { @@ -234,7 +262,7 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi releaseAllCursors(); } - /// The base worker-thread entry point must never run here ({@link #start()} is a no-op). Fail loud if it ever does. + /// The base worker-thread entry point must never run here ([#start()] is a no-op). Fail loud if it ever does. @Override protected void processSegments() { throw new IllegalStateException( @@ -329,12 +357,12 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi /// Activates not-yet-active cursors whose min/max bound can still reach the merge frontier. /// - /// Cursors are visited in min/max order and {@code _nextToActivate} advances only on a real activation, so a - /// {@code break} defers the current cursor to a later call with a risen frontier rather than skipping it. Deferring + /// Cursors are visited in min/max order and `_nextToActivate` advances only on a real activation, so a + /// `break` defers the current cursor to a later call with a risen frontier rather than skipping it. Deferring /// is safe because the first order-by column is the primary sort key: a cursor whose range starts past the frontier /// holds no row sorting before it. With no frontier known yet, activation is forced. Under a single order-by /// expression a cursor whose bound merely *ties* the frontier is deferred too, which is what keeps a long run of - /// equal segment minima from activating every segment at once; see {@link #sortsBeyond}. + /// equal segment minima from activating every segment at once; see [#sortsBeyond]. /// /// The frontier is the row about to be emitted. Understating it defers a cursor that could have supplied that row /// and the merge emits out of order, where overstating it only activates a segment early. The leader is retained @@ -366,11 +394,11 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi /// be compared, so this reports `false` and the caller activates. Tests column 0 only, as the pruning bound always /// has, rather than the full-row comparator -- which would put every order-by column on the per-row path. /// - /// Under a single order-by expression ({@code _deferTiedCursors}) a *tie* also defers, which only postpones a - /// cursor: {@code _nextToActivate} does not advance, and the caller force-activates once heap and leader drain. - /// Since {@code bound} bounds every row in the segment, a deferred cursor holds no row sorting strictly before an + /// Under a single order-by expression (`_deferTiedCursors`) a *tie* also defers, which only postpones a + /// cursor: `_nextToActivate` does not advance, and the caller force-activates once heap and leader drain. + /// Since `bound` bounds every row in the segment, a deferred cursor holds no row sorting strictly before an /// emitted one -- only ties, interchangeable when column 0 is the whole sort key. Note this is *not* - /// {@link MinMaxValueBasedSelectionOrderByCombineOperator}'s justification: its bound is a complete top-K's k-th + /// [MinMaxValueBasedSelectionOrderByCombineOperator]'s justification: its bound is a complete top-K's k-th /// row, this one is the live frontier. It changes which tied rows are returned, already documented as arbitrary. private boolean sortsBeyond(Comparable bound, @Nullable SegmentCursor leader, @Nullable SegmentCursor top) { Object leaderHead = leader != null ? leader.currentHead()[0] : null; @@ -404,7 +432,7 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi } } - /// Returns the accumulated output rows as a sorted {@link SelectionResultsBlock} and resets the output buffer. The + /// Returns the accumulated output rows as a sorted [SelectionResultsBlock] and resets the output buffer. The /// block carries the comparator so the broker-side n-way reduce stays correct. Execution stats are not attached /// here: they go on the terminal metadata block. private BaseResultsBlock flushDataBlock() { @@ -453,9 +481,9 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi } /// 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 + /// loop until the operator returns `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}). + /// [StreamingSelectionOrderByCombineOperator]). private class SegmentCursor { private final Operator<BaseResultsBlock> _operator; @Nullable @@ -478,14 +506,14 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi _maxValue = maxValue; } - /// Returns the current head row to be merged next, or {@code null} if not activated or exhausted. + /// Returns the current head row to be merged next, or `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 + /// After this call [#currentHead()] returns the first row, or `null` if the segment contributes /// nothing. void activate() { acquireSegment(); @@ -495,7 +523,7 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi } /// 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}. + /// when the cursor is exhausted; afterwards [#currentHead()] returns `null`. void advance() { _pos++; if (_pos < _rows.size()) { @@ -510,14 +538,15 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi } /// 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 + /// sight. Returns `false` when the operator is exhausted (no more rows). Streaming-backed operators emit + /// one run/block per call and `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 + /// [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. + /// differing width under one schema would corrupt the result rather than fail. The cursor is exhausted rather than + /// skipping the one block, which drops the same unit the merger does: there, a block is a segment's whole result. private boolean pullBlock() { while (true) { SelectionResultsBlock block = nextBlock(); @@ -562,8 +591,8 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi } } - /// Returns whether the underlying child operator is the lazy {@link StreamingSelectionOrderByOperator}. Must be - /// called after the first {@link Operator#nextBlock()}, which is what materializes a wrapped child. + /// Returns whether the underlying child operator is the lazy [StreamingSelectionOrderByOperator]. Must be + /// called after the first [Operator#nextBlock()], which is what materializes a wrapped child. private boolean isStreamingChild() { Operator underlying = _operator; if (_operator instanceof AcquireReleaseColumnsSegmentOperator) { @@ -583,7 +612,7 @@ public class StreamingSelectionOrderByCombineOperator extends BaseStreamingCombi _acquired = true; } - /// Releases the segment if still held. Idempotent: safe to call from {@link #exhaust()} and combine cleanup. + /// Releases the segment if still held. Idempotent: safe to call from [#exhaust()] and combine cleanup. private void release() { if (_acquired) { if (_operator instanceof AcquireReleaseColumnsSegmentOperator) { 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 index 10dbd97ccde..58f5f37ef7b 100644 --- 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 @@ -61,32 +61,32 @@ 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 +/// Unlike [SelectionOrderByOperator] (which materializes the segment's whole top-K in a single block) this +/// operator emits one globally-sorted [SelectionResultsBlock] per [#getNextBlock()] call and returns +/// `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))}). +/// the query order (the caller must guarantee `projectOperator.isCompatibleWith(DocIdOrder.fromAsc(asc))`). /// -/// Every instance emits at least one block before {@code null}: a segment matching no rows emits a single empty block -/// carrying the {@link DataSchema}, so consumers never have to reconstruct a schema the segment already knows. +/// Every instance emits at least one block before `null`: a segment matching no rows emits a single empty block +/// carrying the [DataSchema], so consumers never have to reconstruct a schema the segment already knows. /// /// It runs in one of two emission modes: /// -/// - **No tail to sort** ({@code numSortedExpressions == numOrderByExpressions}, e.g. {@code ORDER BY sorted}): +/// - **No tail to sort** (`numSortedExpressions == numOrderByExpressions`, e.g. `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}): +/// (trimmed to the remaining `limit + offset` budget). +/// - **Tail to sort** (`numSortedExpressions < numOrderByExpressions`, e.g. +/// `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). +/// `limit + offset` rows by the full comparator, and emits them sorted. This bounds the in-memory run buffer to +/// `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 +/// Like [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 +/// This operator is stateful across [#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"; @@ -231,7 +231,7 @@ public class StreamingSelectionOrderByOperator extends BaseOperator<SelectionRes } /// 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, and otherwise a + /// trimmed to the remaining `limit + offset` budget. Returns `null` when exhausted, and otherwise a /// non-empty list -- empty project blocks are skipped and the row budget is checked up front, so "no rows" always /// means "no more rows". @Nullable @@ -273,7 +273,7 @@ public class StreamingSelectionOrderByOperator extends BaseOperator<SelectionRes } /// 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 + /// `limit + offset` rows by the full comparator, and return them sorted. Returns `null` when /// exhausted. @Nullable private List<Object[]> nextRun() { @@ -314,8 +314,8 @@ public class StreamingSelectionOrderByOperator extends BaseOperator<SelectionRes } /// 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 + /// `Object[_numExpressions]`. For two-phase the document id is stashed at index + /// `_numOrderByExpressions` (overwritten in the second pass). Returns `null` when the project operator /// is exhausted. @Nullable private Object[] nextRow() { @@ -350,7 +350,7 @@ public class StreamingSelectionOrderByOperator extends BaseOperator<SelectionRes return materializeRow(_currentFetcher, _currentDocIds, _currentNullBitmaps, rowId); } - /// Pulls the next project block carrying documents, skipping any that carry none. Returns {@code null} only when + /// Pulls the next project block carrying documents, skipping any that carry none. Returns `null` only when /// the project operator is exhausted; an empty block means "nothing in this batch", not "end of segment", and /// treating one as exhaustion would truncate the scan. /// @@ -410,7 +410,7 @@ public class StreamingSelectionOrderByOperator extends BaseOperator<SelectionRes /// 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()}. + /// instances, mirroring [SelectionOrderByOperator#computePartiallyOrdered()]. private void fetchNonOrderByColumns(List<Object[]> rows) { int numRows = rows.size(); RoaringBitmap docIds = new RoaringBitmap(); @@ -478,9 +478,9 @@ public class StreamingSelectionOrderByOperator extends BaseOperator<SelectionRes } /// Resolves a non-order-by expression's result type without building the phase-2 pipeline, for the zero-match case - /// where there is nothing to fetch. This reproduces {@link TransformOperator#getResultColumnContext} exactly -- that + /// where there is nothing to fetch. This reproduces [TransformOperator#getResultColumnContext] exactly -- that /// method resolves against its project operator's source column contexts, which for phase 2 are - /// {@link #_phase2DataSourceMap} -- so the schema is identical to the one the fetch path builds, at metadata cost + /// [#_phase2DataSourceMap] -- so the schema is identical to the one the fetch path builds, at metadata cost /// only. Worth the indirection because a selective filter can leave many segments of a table matching nothing, and /// each would otherwise stand up a full projection and transform operator to read types it already knows. private ColumnContext resolveResultColumnContext(ExpressionContext expression) { @@ -558,8 +558,8 @@ public class StreamingSelectionOrderByOperator extends BaseOperator<SelectionRes 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 + /// Reports scan cost to the shared [QueryScanCostContext] so scan-based query killing + /// ([org.apache.pinot.core.common.Operator#nextBlock()] -> `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(); 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 e0c8c541184..f00ffb4aa18 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 @@ -116,6 +116,9 @@ public class InstancePlanMakerImplV2 implements PlanMaker { private int _minSegmentGroupTrimSize = Server.DEFAULT_QUERY_EXECUTOR_MIN_SEGMENT_GROUP_TRIM_SIZE; private int _minServerGroupTrimSize = Server.DEFAULT_QUERY_EXECUTOR_MIN_SERVER_GROUP_TRIM_SIZE; private int _groupByTrimThreshold = Server.DEFAULT_QUERY_EXECUTOR_GROUPBY_TRIM_THRESHOLD; + // Server-wide defaults for the streaming selection ORDER BY merge; the query options override them + private double _sortedSelectionMergeAutoMinSortedRatio = Server.DEFAULT_SORTED_SELECTION_MERGE_AUTO_MIN_SORTED_RATIO; + private int _sortedSelectionMergeBlockSize = Server.DEFAULT_SORTED_SELECTION_MERGE_BLOCK_SIZE; @Override public void init(PinotConfiguration queryExecutorConfig) { @@ -146,11 +149,24 @@ public class InstancePlanMakerImplV2 implements PlanMaker { Server.DEFAULT_QUERY_EXECUTOR_GROUPBY_TRIM_THRESHOLD); Preconditions.checkState(_groupByTrimThreshold > 0, "Invalid configurable: groupByTrimThreshold: %d must be positive", _groupByTrimThreshold); + _sortedSelectionMergeAutoMinSortedRatio = + queryExecutorConfig.getProperty(Server.SORTED_SELECTION_MERGE_AUTO_MIN_SORTED_RATIO, + Server.DEFAULT_SORTED_SELECTION_MERGE_AUTO_MIN_SORTED_RATIO); + Preconditions.checkState( + _sortedSelectionMergeAutoMinSortedRatio >= 0 && _sortedSelectionMergeAutoMinSortedRatio <= 1, + "Invalid configuration: sortedSelectionMergeAutoMinSortedRatio: %s must be in [0, 1]", + _sortedSelectionMergeAutoMinSortedRatio); + _sortedSelectionMergeBlockSize = queryExecutorConfig.getProperty(Server.SORTED_SELECTION_MERGE_BLOCK_SIZE, + Server.DEFAULT_SORTED_SELECTION_MERGE_BLOCK_SIZE); + Preconditions.checkState(_sortedSelectionMergeBlockSize > 0, + "Invalid configuration: sortedSelectionMergeBlockSize: %s must be positive", _sortedSelectionMergeBlockSize); LOGGER.info("Initialized plan maker with maxExecutionThreads: {}, defaultExecutionThreads: {}, " + "maxInitialResultHolderCapacity: {}, numGroupsLimit: {}, minSegmentGroupTrimSize: {}, " - + "minServerGroupTrimSize: {}, groupByTrimThreshold: {}", + + "minServerGroupTrimSize: {}, groupByTrimThreshold: {}, sortedSelectionMergeAutoMinSortedRatio: {}, " + + "sortedSelectionMergeBlockSize: {}", _maxExecutionThreads, _defaultExecutionThreads, _maxInitialResultHolderCapacity, _numGroupsLimit, - _minSegmentGroupTrimSize, _minServerGroupTrimSize, _groupByTrimThreshold); + _minSegmentGroupTrimSize, _minServerGroupTrimSize, _groupByTrimThreshold, + _sortedSelectionMergeAutoMinSortedRatio, _sortedSelectionMergeBlockSize); } @VisibleForTesting @@ -291,14 +307,14 @@ public class InstancePlanMakerImplV2 implements PlanMaker { } Double sortedSelectionMergeAutoMinSortedRatio = QueryOptionsUtils.getSortedSelectionMergeAutoMinSortedRatio(queryOptions); - if (sortedSelectionMergeAutoMinSortedRatio != null) { - queryContext.setSortedSelectionMergeAutoMinSortedRatio(sortedSelectionMergeAutoMinSortedRatio); - } + queryContext.setSortedSelectionMergeAutoMinSortedRatio(sortedSelectionMergeAutoMinSortedRatio != null + ? sortedSelectionMergeAutoMinSortedRatio + : _sortedSelectionMergeAutoMinSortedRatio); Integer sortedSelectionMergeBlockSize = QueryOptionsUtils.getSortedSelectionMergeBlockSize(queryOptions); - if (sortedSelectionMergeBlockSize != null) { - queryContext.setSortedSelectionMergeBlockSize(sortedSelectionMergeBlockSize); - } + queryContext.setSortedSelectionMergeBlockSize(sortedSelectionMergeBlockSize != null + ? sortedSelectionMergeBlockSize + : _sortedSelectionMergeBlockSize); } // Set group-by query options 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 384f800c8b2..7a5741826e6 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,7 +44,6 @@ 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; import org.apache.pinot.spi.utils.CommonConstants.Server.SortedSelectionMergeMode; import org.slf4j.Logger; @@ -148,7 +147,7 @@ public class QueryContext { private double _sortedSelectionMergeAutoMinSortedRatio = Server.DEFAULT_SORTED_SELECTION_MERGE_AUTO_MIN_SORTED_RATIO; /// Output block size (rows) for the streaming selection ORDER BY combine - private int _sortedSelectionMergeBlockSize = Broker.DEFAULT_SORTED_SELECTION_MERGE_BLOCK_SIZE; + private int _sortedSelectionMergeBlockSize = Server.DEFAULT_SORTED_SELECTION_MERGE_BLOCK_SIZE; // Guards the one-time warning in isSortedSelectionMergeEnabled() so that an unresolved AUTO does not log once per // segment/combine call for the same query. private volatile boolean _unresolvedSortedSelectionMergeModeWarned; 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 da476ee0dbd..786bebaf314 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 @@ -179,9 +179,10 @@ public class CombineSlowOperatorsTest { /// 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 + /// timeout **before** activating any child - asserted via `_operationInProgress`, which also keeps the test + /// from passing vacuously if the timeout came from somewhere else. A regression would drive a [SlowOperator], which + /// sleeps for an hour, so the test timeout turns that hang into a failure. + @Test(timeOut = 30_000) public void testStreamingSelectionOrderByCombineOperatorHonorsDeadline() { // A real latch, never awaited: this test asserts that no child is driven at all, so there is nothing to wait // for. It is passed only so every getOperators() call site in this class looks alike; SlowOperator null-guards @@ -221,7 +222,7 @@ public class CombineSlowOperatorsTest { } /// 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 + /// asserts the operator surfaces an [ExceptionResultsBlock]. When `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. diff --git a/pinot-core/src/test/java/org/apache/pinot/core/operator/combine/StreamingSelectionOrderByCombineOperatorTest.java b/pinot-core/src/test/java/org/apache/pinot/core/operator/combine/StreamingSelectionOrderByCombineOperatorTest.java index ae269cc1e76..606705708cc 100644 --- a/pinot-core/src/test/java/org/apache/pinot/core/operator/combine/StreamingSelectionOrderByCombineOperatorTest.java +++ b/pinot-core/src/test/java/org/apache/pinot/core/operator/combine/StreamingSelectionOrderByCombineOperatorTest.java @@ -23,18 +23,25 @@ import java.io.IOException; import java.util.ArrayList; import java.util.Comparator; import java.util.List; +import java.util.Map; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.stream.Collectors; +import javax.annotation.Nullable; import org.apache.commons.io.FileUtils; +import org.apache.pinot.common.proto.Plan.ExplainNode.AttributeValue; 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.ExceptionResultsBlock; 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.plan.CombinePlanNode; +import org.apache.pinot.core.plan.ExplainInfo; import org.apache.pinot.core.plan.PlanNode; import org.apache.pinot.core.plan.maker.InstancePlanMakerImplV2; import org.apache.pinot.core.plan.maker.PlanMaker; @@ -54,6 +61,8 @@ 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.exception.QueryErrorCode; +import org.apache.pinot.spi.exception.QueryErrorMessage; import org.apache.pinot.spi.utils.CommonConstants.Server; import org.apache.pinot.spi.utils.CommonConstants.Server.SortedSelectionMergeMode; import org.apache.pinot.spi.utils.ReadMode; @@ -71,22 +80,22 @@ import static org.testng.Assert.assertTrue; import static org.testng.Assert.expectThrows; -/// Combine-level tests for {@link StreamingSelectionOrderByCombineOperator} (step-3 operator) and its wiring into -/// {@link CombinePlanNode#getCombineOperator()} (step-4). +/// Combine-level tests for [StreamingSelectionOrderByCombineOperator] (step-3 operator) and its wiring into +/// [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 sortedSelectionMergeMode=ON} (asserting the new -/// operator was actually selected) and once with the hint off (asserting the {@code MinMax} operator was selected) - +/// The streaming combine must return the same globally-sorted top-K rows as the default +/// [MinMaxValueBasedSelectionOrderByCombineOperator], only (in streaming mode) spread across several bounded +/// blocks. Each functional test therefore asserts **streaming-vs-non-streaming parity**: it runs the identical query +/// twice over the same in-memory segments - once with `sortedSelectionMergeMode=ON` (asserting the new +/// operator was actually selected) and once with the hint off (asserting the `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 +/// 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} +/// `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 +/// (`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 = @@ -101,6 +110,10 @@ public class StreamingSelectionOrderByCombineOperatorTest { /// 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"; + /// Takes the LIMIT. `tailCol` descends within each `_sortedSegments` segment, so it is not physically sorted and + /// each streaming child takes the run path: one block per distinct sortedCol value, which is one row per block. + private static final String ONE_ROW_PER_BLOCK_QUERY = + "SELECT sortedCol, tailCol, valCol FROM testTable ORDER BY sortedCol, tailCol LIMIT %d"; /// 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; @@ -452,7 +465,7 @@ public class StreamingSelectionOrderByCombineOperatorTest { /// Correctness under single-column ties with an OFFSET: LIMIT 10 OFFSET 20 straddles the same sortedCol=0/1 /// boundary as [#testSingleColumnTieDeferralPreservesOrderByValuesAcrossLimitStraddle] (limit + offset = 30), but /// exercises it with a nonzero offset. As [#testLimitOffsetParity] documents, the server (and this combine) retains - /// {@code limit + offset} rows -- the offset is trimmed by the broker afterwards -- so 30, not 10, rows come back + /// `limit + offset` rows -- the offset is trimmed by the broker afterwards -- so 30, not 10, rows come back /// here too; what differs from the straddle test is only the query shape, not the row count. @Test public void testSingleColumnTieDeferralPreservesOrderByValuesAcrossOffset() { @@ -617,6 +630,9 @@ public class StreamingSelectionOrderByCombineOperatorTest { + 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); + // The pruning is visible in the stats only as the gap between these two: see the class javadoc. + assertEquals(result._numSegmentsProcessed, NUM_SEGMENTS, "numSegmentsProcessed counts every segment"); + assertEquals(result._numSegmentsMatched, 1, "numSegmentsMatched must exclude never-activated segments"); } /// A cursor activated *late* can hold a smaller row than the retained leader. Two things must be right for it to @@ -756,6 +772,219 @@ public class StreamingSelectionOrderByCombineOperatorTest { expectThrows(IllegalStateException.class, streamingCombine::processSegments); } + // Segment acquire/release lifecycle. The merge acquires a segment when it activates the segment's cursor, so every + // way out of the merge has to release whatever is still held; a missed release pins the segment for the lifetime of + // the server. `_sortedSegments` segment i covers sortedCol [i * 50, i * 50 + 100). + + @Test + public void testFullDrainReleasesEverySegment() { + QueryContext queryContext = + hintedContext("SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol LIMIT 100000"); + List<InstrumentedSegmentOperator> children = instrument(_sortedSegments, queryContext); + List<BaseResultsBlock> blocks = drain(combineOver(children, queryContext)); + assertEquals(rowsOf(blocks).size(), NUM_SEGMENTS * NUM_RECORDS_PER_SEGMENT); + assertEquals(numAcquired(children), NUM_SEGMENTS, "A full drain must activate every segment"); + assertEveryAcquireReleased(children); + } + + @Test + public void testLimitReachedReleasesSegmentsStillHeld() { + // 60 rows end around sortedCol 55, where segments 0 and 1 are both active and neither is drained, so finish() + // is the only thing that can release them. Segment 2 starts at 100 and must never be acquired. + QueryContext queryContext = + hintedContext("SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol LIMIT 60"); + queryContext.setSortedSelectionMergeBlockSize(1000); + List<InstrumentedSegmentOperator> children = instrument(_sortedSegments, queryContext); + List<BaseResultsBlock> blocks = drain(combineOver(children, queryContext)); + assertEquals(rowsOf(blocks).size(), 60); + assertEquals(numAcquired(children), 2, "Only the two segments overlapping the first 60 rows may be acquired"); + assertEveryAcquireReleased(children); + } + + @Test + public void testStopReleasesSegmentsHeldMidMerge() { + QueryContext queryContext = + hintedContext("SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol LIMIT 200"); + queryContext.setSortedSelectionMergeBlockSize(60); + List<InstrumentedSegmentOperator> children = instrument(_sortedSegments, queryContext); + StreamingSelectionOrderByCombineOperator combine = combineOver(children, queryContext); + + BaseResultsBlock first = combine.nextBlock(); + assertTrue(first instanceof SelectionResultsBlock, "Expected a data block, got: " + first.getClass()); + assertTrue(numHeld(children) > 0, "Segments must be held mid-merge, else the test is vacuous"); + + combine.stop(); + assertEveryAcquireReleased(children); + // A driver may stop more than once; release must not run again for a segment already released. + combine.stop(); + assertEveryAcquireReleased(children); + assertTrue(combine.nextBlock() instanceof MetadataResultsBlock, "A stopped merge must not resume"); + assertEveryAcquireReleased(children); + } + + @Test + public void testChildExceptionMidMergeReleasesEverySegment() { + // Segment 0 fails on its 60th block (sortedCol 59), after segment 1 (which starts at 50) has been activated, so a + // segment other than the failing one is also held when the exception unwinds the merge. + QueryContext queryContext = hintedContext(String.format(ONE_ROW_PER_BLOCK_QUERY, 200)); + queryContext.setSortedSelectionMergeBlockSize(1000); + List<InstrumentedSegmentOperator> children = instrument(_sortedSegments, queryContext); + children.get(0)._interceptor = (blockNumber, block) -> { + if (blockNumber == 60) { + throw new IllegalStateException("Injected child failure"); + } + return block; + }; + BaseResultsBlock result = combineOver(children, queryContext).nextBlock(); + assertTrue(result instanceof ExceptionResultsBlock, "Expected an exception block, got: " + result.getClass()); + assertEquals(numAcquired(children), 2, "Segments 0 and 1 must both be held when the failure hits"); + assertEveryAcquireReleased(children); + } + + @Test + public void testChildErrorMidMergeReleasesEverySegment() { + // An Error is rethrown rather than turned into an exception block, so single-stage execution has no stop() after + // it: the merge's own Throwable handler is the only release. + QueryContext queryContext = hintedContext(String.format(ONE_ROW_PER_BLOCK_QUERY, 200)); + queryContext.setSortedSelectionMergeBlockSize(1000); + List<InstrumentedSegmentOperator> children = instrument(_sortedSegments, queryContext); + children.get(0)._interceptor = (blockNumber, block) -> { + if (blockNumber == 60) { + throw new OutOfMemoryError("Injected child error"); + } + return block; + }; + StreamingSelectionOrderByCombineOperator combine = combineOver(children, queryContext); + expectThrows(OutOfMemoryError.class, combine::nextBlock); + assertEquals(numAcquired(children), 2, "Segments 0 and 1 must both be held when the error hits"); + assertEveryAcquireReleased(children); + } + + /// A block whose schema disagrees with the merge's is dropped together with the rest of its segment, and reported + /// once as a MERGE_RESPONSE error however many segments disagree the same way. A segment's schema does not change + /// between its own blocks in practice (a reload swaps the whole segment); diverging on a middle block here is what + /// separates "drop the rest of the segment" from "drop one block". + @Test + public void testSchemaMismatchDropsTheRestOfTheSegmentAndIsReportedOnce() { + String query = String.format(ONE_ROW_PER_BLOCK_QUERY, 100000); + QueryContext queryContext = hintedContext(query); + queryContext.setSortedSelectionMergeBlockSize(7); + List<InstrumentedSegmentOperator> children = instrument(_sortedSegments, queryContext); + int divergingBlock = 11; + for (int segment : List.of(1, 2)) { + children.get(segment)._interceptor = (blockNumber, block) -> blockNumber == divergingBlock + ? withMismatchedSchema((SelectionResultsBlock) block) + : block; + } + + List<BaseResultsBlock> blocks = drain(combineOver(children, queryContext)); + assertTrue(blocks.get(blocks.size() - 1) instanceof MetadataResultsBlock, + "A schema mismatch must not fail the query, got: " + blocks.get(blocks.size() - 1).getClass()); + List<Object[]> rows = rowsOf(blocks); + assertSorted(rows, orderByComparator(query, false)); + + // One row per block, so each diverging segment contributes exactly the rows of the blocks before the mismatch. + int valColIndex = List.of(blocks.get(0).getDataSchema().getColumnNames()).indexOf(VAL_COL); + int[] rowsPerSegment = new int[NUM_SEGMENTS]; + for (Object[] row : rows) { + int valCol = (Integer) row[valColIndex]; + int segment = valCol / 1_000_000; + rowsPerSegment[segment]++; + if (segment == 1 || segment == 2) { + assertTrue(valCol % 1_000_000 < divergingBlock - 1, + "Segment " + segment + " emitted a row from at or after the mismatched block: " + valCol); + } + } + for (int segment = 0; segment < NUM_SEGMENTS; segment++) { + int expected = segment == 1 || segment == 2 ? divergingBlock - 1 : NUM_RECORDS_PER_SEGMENT; + assertEquals(rowsPerSegment[segment], expected, "Unexpected row count for segment " + segment); + } + for (int segment : List.of(1, 2)) { + assertEquals(children.get(segment)._numBlocks, divergingBlock, + "Segment " + segment + " must not be read past the mismatched block"); + } + assertEveryAcquireReleased(children); + + // The same mismatch in two segments produces one message, attached to one block rather than re-attached to every + // block streamed after it. + List<QueryErrorMessage> errors = new ArrayList<>(); + for (BaseResultsBlock block : blocks) { + if (block.getErrorMessages() != null) { + errors.addAll(block.getErrorMessages()); + } + } + assertEquals(errors.size(), 1, "Expected exactly one mismatch error, got: " + errors); + assertEquals(errors.get(0).getErrCode(), QueryErrorCode.MERGE_RESPONSE); + } + + /// Returns `block` with its last column retyped from INT to LONG, the shape a newly reloaded segment takes when a + /// column's type changes. + private static SelectionResultsBlock withMismatchedSchema(SelectionResultsBlock block) { + DataSchema dataSchema = block.getDataSchema(); + ColumnDataType[] columnDataTypes = dataSchema.getColumnDataTypes().clone(); + int last = columnDataTypes.length - 1; + assertEquals(columnDataTypes[last], ColumnDataType.INT); + columnDataTypes[last] = ColumnDataType.LONG; + return new SelectionResultsBlock(new DataSchema(dataSchema.getColumnNames(), columnDataTypes), block.getRows(), + block.getComparator(), block.getQueryContext()); + } + + /// Plans one [InstrumentedSegmentOperator] per segment, in segment order, the way the prefetch path wraps the real + /// streaming leaf. + private static List<InstrumentedSegmentOperator> instrument(List<IndexSegment> segments, + QueryContext queryContext) { + List<InstrumentedSegmentOperator> children = new ArrayList<>(segments.size()); + for (IndexSegment segment : segments) { + children.add(new InstrumentedSegmentOperator( + PLAN_MAKER.makeStreamingSegmentPlanNode(new SegmentContext(segment), queryContext), segment)); + } + return children; + } + + private static StreamingSelectionOrderByCombineOperator combineOver(List<InstrumentedSegmentOperator> children, + QueryContext queryContext) { + return new StreamingSelectionOrderByCombineOperator(new ArrayList<>(children), queryContext, EXECUTOR); + } + + /// Drives `combine` until it returns something other than a data block, and returns every block, that one last. + private static List<BaseResultsBlock> drain(Operator<?> combine) { + List<BaseResultsBlock> blocks = new ArrayList<>(); + while (true) { + BaseResultsBlock block = (BaseResultsBlock) combine.nextBlock(); + blocks.add(block); + if (!(block instanceof SelectionResultsBlock)) { + return blocks; + } + assertTrue(blocks.size() < 1_000_000, "Streaming combine did not terminate"); + } + } + + private static List<Object[]> rowsOf(List<BaseResultsBlock> blocks) { + List<Object[]> rows = new ArrayList<>(); + for (BaseResultsBlock block : blocks) { + if (block instanceof SelectionResultsBlock) { + rows.addAll(((SelectionResultsBlock) block).getRows()); + } + } + return rows; + } + + private static int numAcquired(List<InstrumentedSegmentOperator> children) { + return children.stream().mapToInt(child -> child._numAcquires).sum(); + } + + private static int numHeld(List<InstrumentedSegmentOperator> children) { + return children.stream().mapToInt(child -> child._numAcquires - child._numReleases).sum(); + } + + private static void assertEveryAcquireReleased(List<InstrumentedSegmentOperator> children) { + for (int i = 0; i < children.size(); i++) { + InstrumentedSegmentOperator child = children.get(i); + assertTrue(child._numAcquires <= 1, "Segment " + i + " acquired " + child._numAcquires + " times"); + assertEquals(child._numReleases, child._numAcquires, "Segment " + i + " acquires and releases differ"); + } + } + @Test public void testHintOffSelectsMinMaxOperator() { // Default behavior is unchanged when the hint is off: the classic MinMax operator is still selected. @@ -815,7 +1044,65 @@ public class StreamingSelectionOrderByCombineOperatorTest { + combineOperator.getClass().getSimpleName()); } - /// Asserts streaming-vs-blocking parity for {@code query}. Runs the MinMax combine on the blocking path (hint off) + @Test + public void testExplainAttributes() { + @Language("sql") String query = "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol LIMIT 50"; + QueryContext queryContext = hintedContext(query); + queryContext.setSortedSelectionMergeBlockSize(7); + Operator<?> combineOperator = planCombineOperator(_sortedSegments, queryContext, true); + ExplainInfo explainInfo = combineOperator.getExplainInfo(); + assertEquals(explainInfo.getTitle(), "CombineSelectOrderbyStreaming"); + Map<String, AttributeValue> attributes = explainInfo.getAttributes(); + assertEquals(attributes.get("blockSize").getLong(), 7); + // Query-wide settings must not be summed when the broker merges this node across workers. + assertEquals(attributes.get("blockSize").getMergeType(), AttributeValue.MergeType.IDEMPOTENT); + assertTrue(attributes.get("frontierPruning").getBool()); + assertFalse(attributes.get("deferTiedCursors").getBool(), "Two order-by expressions must not defer ties"); + assertEquals(attributes.get("numSegments").getLong(), NUM_SEGMENTS); + assertEquals(attributes.get("numSortedSegments").getLong(), NUM_SEGMENTS); + + // Explaining must not start the merge: the same operator still returns the full, correct result afterwards. + List<Object[]> rows = new ArrayList<>(); + while (true) { + BaseResultsBlock block = (BaseResultsBlock) combineOperator.nextBlock(); + if (block instanceof MetadataResultsBlock) { + break; + } + rows.addAll(((SelectionResultsBlock) block).getRows()); + } + assertMultisetEquals(rows, run(_sortedSegments, query, false, false, false, 0)._rows); + } + + @Test + public void testExplainAttributesReflectTheQueryShape() { + assertTrue(explainAttributes(_sortedSegments, "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol LIMIT 5", + false).get("deferTiedCursors").getBool(), "A single order-by expression defers ties"); + assertFalse(explainAttributes(_sortedSegments, + "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol LIMIT 5", true) + .get("frontierPruning").getBool(), "Null handling disables frontier pruning"); + assertEquals(explainAttributes(_mixedSegments, + "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol LIMIT 5", false) + .get("numSortedSegments").getLong(), 2); + assertEquals(explainAttributes(_unsortedSegments, + "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol LIMIT 5", false) + .get("numSortedSegments").getLong(), 0); + + Map<String, AttributeValue> expression = explainAttributes(_sortedSegments, + "SELECT sortedCol, valCol FROM testTable ORDER BY ADD(sortedCol, 1), valCol LIMIT 5", false); + assertFalse(expression.get("frontierPruning").getBool(), "An expression leading order-by has no min/max to prune"); + assertFalse(expression.containsKey("numSortedSegments"), "Sortedness is undefined for an expression"); + } + + private static Map<String, AttributeValue> explainAttributes(List<IndexSegment> segments, + @Language("sql") String query, boolean nullHandling) { + QueryContext queryContext = hintedContext(query); + queryContext.setNullHandlingEnabled(nullHandling); + Operator<?> combineOperator = planCombineOperator(segments, queryContext, true); + assertTrue(combineOperator instanceof StreamingSelectionOrderByCombineOperator); + return combineOperator.getExplainInfo().getAttributes(); + } + + /// Asserts streaming-vs-blocking parity for `query`. Runs the MinMax combine on the blocking path (hint off) /// as the reference, then runs the streaming combine on the streaming path, asserting it selects the streaming /// operator and produces rows that are sorted by the order-by comparator and equal the MinMax rows as a multiset, /// plus the bounded-flush invariants. A small block size forces several bounded data blocks before the metadata @@ -852,7 +1139,7 @@ public class StreamingSelectionOrderByCombineOperatorTest { } } - /// Plans the combine operator that would run {@code queryContext} over {@code segments}, without driving it. Use + /// Plans the combine operator that would run `queryContext` over `segments`, without driving it. Use /// this for gate assertions: an operator that is not the streaming combine (the fallbacks) does not terminate with /// a metadata block, so it cannot be driven by the streaming loop in [#run]. private static Operator<?> planCombineOperator(List<IndexSegment> segments, QueryContext queryContext, @@ -898,17 +1185,17 @@ public class StreamingSelectionOrderByCombineOperatorTest { @Test public void testAutoHonoursTheMinSortedRatioThreshold() { // _mixedSegments is 2 sorted of 4, so a ratio of exactly 0.5 must pass and anything above it must fail. This pins - // the comparison as >= rather than >, and pins that the threshold is read from the query context. + // the comparison as >= rather than >, and pins that the threshold is read from the query option. String query = "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol, valCol LIMIT 50"; - QueryContext atThreshold = modeContext(query, SortedSelectionMergeMode.AUTO); - atThreshold.setSortedSelectionMergeAutoMinSortedRatio(0.5); + QueryContext atThreshold = + modeContext("SET sortedSelectionMergeAutoMinSortedRatio=0.5; " + query, SortedSelectionMergeMode.AUTO); makeStreamingInstancePlan(_mixedSegments, atThreshold); assertEquals(atThreshold.getSortedSelectionMergeMode(), SortedSelectionMergeMode.ON, "A sorted ratio equal to the threshold must select the streaming merge"); - QueryContext aboveThreshold = modeContext(query, SortedSelectionMergeMode.AUTO); - aboveThreshold.setSortedSelectionMergeAutoMinSortedRatio(0.75); + QueryContext aboveThreshold = + modeContext("SET sortedSelectionMergeAutoMinSortedRatio=0.75; " + query, SortedSelectionMergeMode.AUTO); makeStreamingInstancePlan(_mixedSegments, aboveThreshold); assertEquals(aboveThreshold.getSortedSelectionMergeMode(), SortedSelectionMergeMode.OFF, "A sorted ratio below the threshold must keep the MinMax combine"); @@ -1018,14 +1305,14 @@ public class StreamingSelectionOrderByCombineOperatorTest { String query = "SET allowReverseOrder=true; SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol DESC, " + "valCol DESC LIMIT 50"; - QueryContext atThreshold = modeContext(query, SortedSelectionMergeMode.AUTO); - atThreshold.setSortedSelectionMergeAutoMinSortedRatio(0.5); + QueryContext atThreshold = + modeContext("SET sortedSelectionMergeAutoMinSortedRatio=0.5; " + query, SortedSelectionMergeMode.AUTO); makeStreamingInstancePlan(_mixedSegments, atThreshold); assertEquals(atThreshold.getSortedSelectionMergeMode(), SortedSelectionMergeMode.ON, "A DESC sorted ratio equal to the threshold must select the streaming merge"); - QueryContext aboveThreshold = modeContext(query, SortedSelectionMergeMode.AUTO); - aboveThreshold.setSortedSelectionMergeAutoMinSortedRatio(0.75); + QueryContext aboveThreshold = + modeContext("SET sortedSelectionMergeAutoMinSortedRatio=0.75; " + query, SortedSelectionMergeMode.AUTO); makeStreamingInstancePlan(_mixedSegments, aboveThreshold); assertEquals(aboveThreshold.getSortedSelectionMergeMode(), SortedSelectionMergeMode.OFF, "A DESC sorted ratio below the threshold must keep the MinMax combine"); @@ -1069,7 +1356,7 @@ public class StreamingSelectionOrderByCombineOperatorTest { false); } - /// Runs the streaming instance plan for its side effect on {@code queryContext}: resolving AUTO from segment + /// Runs the streaming instance plan for its side effect on `queryContext`: resolving AUTO from segment /// metadata. The plan itself is discarded; the gate tests assert on the resolved mode and on the operator that /// [#planCombineOperator] then selects. private static void makeStreamingInstancePlan(List<IndexSegment> segments, QueryContext queryContext) { @@ -1081,7 +1368,7 @@ public class StreamingSelectionOrderByCombineOperatorTest { }); } - /// Builds a query context for {@code query} with the streaming merge mode set to {@code mode}. + /// Builds a query context for `query` with the streaming merge mode set to `mode`. private static QueryContext modeContext(@Language("sql") String query, SortedSelectionMergeMode mode) { QueryContext queryContext = QueryContextConverterUtils.getQueryContext(query); queryContext.setSortedSelectionMergeMode(mode); @@ -1089,7 +1376,7 @@ public class StreamingSelectionOrderByCombineOperatorTest { return queryContext; } - /// Builds a query context for {@code query} with the streaming merge hint on. + /// Builds a query context for `query` with the streaming merge hint on. private static QueryContext hintedContext(@Language("sql") String query) { QueryContext queryContext = QueryContextConverterUtils.getQueryContext(query); queryContext.setSortedSelectionMergeMode(SortedSelectionMergeMode.ON); @@ -1097,7 +1384,7 @@ public class StreamingSelectionOrderByCombineOperatorTest { return queryContext; } - /// Runs one combine over {@code segments} and collects its rows, blocks, schema and docs-scanned stat. + /// Runs one combine over `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); @@ -1126,6 +1413,7 @@ public class StreamingSelectionOrderByCombineOperatorTest { } result._numDocsScanned = block.getNumDocsScanned(); result._numSegmentsMatched = block.getNumSegmentsMatched(); + result._numSegmentsProcessed = block.getNumSegmentsProcessed(); break; } SelectionResultsBlock dataBlock = (SelectionResultsBlock) block; @@ -1173,7 +1461,7 @@ public class StreamingSelectionOrderByCombineOperatorTest { /// 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. + /// a plain `Arrays.toString` (type-blind) comparison would miss. private static List<String> toCanonical(List<Object[]> rows) { return rows.stream().map(row -> { StringBuilder sb = new StringBuilder("["); @@ -1218,5 +1506,44 @@ public class StreamingSelectionOrderByCombineOperatorTest { private int _numBlocks; private long _numDocsScanned; private int _numSegmentsMatched; + private int _numSegmentsProcessed; + } + + /// A segment child wrapped the way the prefetch path wraps it, counting acquire/release calls and optionally + /// replacing or failing a numbered child block. Acquire and release are counted rather than forwarded, since these + /// immutable test segments have nothing to fetch. + private static class InstrumentedSegmentOperator extends AcquireReleaseColumnsSegmentOperator { + private int _numAcquires; + private int _numReleases; + private int _numBlocks; + @Nullable + private BlockInterceptor _interceptor; + + InstrumentedSegmentOperator(PlanNode planNode, IndexSegment indexSegment) { + super(planNode, indexSegment, null); + } + + @Override + public void acquire() { + _numAcquires++; + } + + @Override + public void release() { + _numReleases++; + } + + @Override + protected BaseResultsBlock getNextBlock() { + BaseResultsBlock block = super.getNextBlock(); + _numBlocks++; + return _interceptor != null ? _interceptor.intercept(_numBlocks, block) : block; + } + } + + @FunctionalInterface + private interface BlockInterceptor { + /// Returns the block to hand the merge in place of `block`, the child's `blockNumber`-th (1-based). + BaseResultsBlock intercept(int blockNumber, BaseResultsBlock block); } } 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 index f77bdc4d4b4..62a83bc7816 100644 --- 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 @@ -73,20 +73,20 @@ import static org.testng.Assert.assertNull; import static org.testng.Assert.assertTrue; -/// Segment-level tests for {@link StreamingSelectionOrderByOperator}. +/// Segment-level tests for [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 sortedSelectionMergeMode=ON}, +/// 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 **stream-vs-materialized parity**: it +/// drives the streaming operator through [SelectionPlanNode] with `sortedSelectionMergeMode=ON`, /// concatenates -/// every {@link Operator#nextBlock()} output until {@code null}, and asserts the concatenation equals the single block -/// the materialized operator ({@link SelectionPartiallyOrderedByLinearOperator} / {@link SelectionOrderByOperator}) +/// every [Operator#nextBlock()] output until `null`, and asserts the concatenation equals the single block +/// the materialized operator ([SelectionPartiallyOrderedByLinearOperator] / [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 +/// 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 `_segment` fixture has a +/// unique sorted column, while the `_dupSegment` / `_largeSegment` fixtures repeat the sorted column but +/// pair it with a unique `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 { @@ -202,7 +202,7 @@ public class StreamingSelectionOrderByOperatorTest { return records; } - /// Appends one run of {@code runSize} rows all sharing {@code sortedValue}, with a tail that descends within the run + /// Appends one run of `runSize` rows all sharing `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++) { @@ -349,7 +349,7 @@ public class StreamingSelectionOrderByOperatorTest { "Expected no nulls in the output when null handling is disabled"); } - /// No-tail mode, single phase: {@link StreamingSelectionOrderByOperator} must skip a zero-document project block + /// No-tail mode, single phase: [StreamingSelectionOrderByOperator] must skip a zero-document project block /// rather than read it as end-of-segment. @Test public void testEmptyProjectBlockDoesNotTruncateSortedScan() { @@ -357,14 +357,14 @@ public class StreamingSelectionOrderByOperatorTest { 1, 0); } - /// Same, two phase: this is also the only path that asks the injected block for {@code getDocIds()}. + /// Same, two phase: this is also the only path that asks the injected block for `getDocIds()`. @Test public void testEmptyProjectBlockDoesNotTruncateTwoPhaseSortedScan() { assertEmptyProjectBlocksAreSkipped(_segment, "SELECT sortedCol, valCol FROM testTable ORDER BY sortedCol LIMIT 25", 1, 2, 1, 0); } - /// Tail mode: the run scan reaches the project operator through {@code nextRow()}, which already tolerated empty + /// Tail mode: the run scan reaches the project operator through `nextRow()`, which already tolerated empty /// blocks. Pins that the two scan paths agree. @Test public void testEmptyProjectBlockDoesNotTruncateRunScan() { @@ -380,7 +380,7 @@ public class StreamingSelectionOrderByOperatorTest { } /// Empty blocks arriving immediately before exhaustion exercise the loop's other exit: the skip must fall through - /// to the project operator's {@code null} rather than spin or emit a phantom row. The limit exceeds the fixture so + /// to the project operator's `null` rather than spin or emit a phantom row. The limit exceeds the fixture so /// the scan actually reaches the end instead of stopping on the row budget. @Test public void testEmptyProjectBlocksBeforeExhaustionEndTheScan() { @@ -388,14 +388,14 @@ public class StreamingSelectionOrderByOperatorTest { 0, 2); } - /// The skip loop sits directly on top of the {@code limit + offset} budget bookkeeping, so cover a non-zero offset. + /// The skip loop sits directly on top of the `limit + offset` budget bookkeeping, so cover a non-zero offset. @Test public void testEmptyProjectBlockIsSkippedWithOffset() { assertEmptyProjectBlocksAreSkipped(_segment, "SELECT sortedCol FROM testTable ORDER BY sortedCol LIMIT 15 OFFSET 10", 1, 2, 1, 0); } - /// {@code nextRow()} rebuilds the phase-1 null bitmaps in the same branch that pulls the next non-empty block, so + /// `nextRow()` rebuilds the phase-1 null bitmaps in the same branch that pulls the next non-empty block, so /// the skip must not desynchronise a bitmap from the block it was built against. Uses the fixture that carries a /// null in the order-by tail column, with null handling on. @Test @@ -489,7 +489,7 @@ public class StreamingSelectionOrderByOperatorTest { } } - /// Each {@link ProjectPlanNode} build looks up one data source per projected column, so a plan that builds the + /// Each [ProjectPlanNode] build looks up one data source per projected column, so a plan that builds the /// project twice looks up more than the materialized plan, which builds it once. @Test public void testDescIncompatibleFallbackBuildsTheProjectOnce() { @@ -504,7 +504,7 @@ public class StreamingSelectionOrderByOperatorTest { "The DESC-incompatible fall-through built the project twice"); } - /// Plans {@code query} in {@code mode} against a segment counting data source lookups into {@code numLookups}, + /// Plans `query` in `mode` against a segment counting data source lookups into `numLookups`, /// asserting it lands on the materialized DESC operator. private Operator<SelectionResultsBlock> planWithCountedDataSourceLookups(@Language("sql") String query, SortedSelectionMergeMode mode, AtomicInteger numLookups) { @@ -525,9 +525,9 @@ public class StreamingSelectionOrderByOperatorTest { } /// Every block the tail-to-sort path emits carries a structurally mutable row list. Defensive rather than a fix for - /// a reachable failure: no consumer of this operator adds to the list today, but {@link SelectionResultsBlock} is a - /// shared type and {@code SelectionOperatorUtils.mergeWithoutOrdering()} adds to the row list of the block it merges - /// into, which a fixed-size {@code Arrays.asList} view would reject. + /// a reachable failure: no consumer of this operator adds to the list today, but [SelectionResultsBlock] is a + /// shared type and `SelectionOperatorUtils.mergeWithoutOrdering()` adds to the row list of the block it merges + /// into, which a fixed-size `Arrays.asList` view would reject. @Test public void testRunPathEmitsGrowableRowLists() { QueryContext queryContext = QueryContextConverterUtils.getQueryContext( @@ -570,7 +570,7 @@ public class StreamingSelectionOrderByOperatorTest { } /// A segment matching no rows must still emit exactly one block before signalling exhaustion: empty, but carrying - /// the same {@link DataSchema} the identical query produces when it does match rows. Consumers therefore never have + /// the same [DataSchema] the identical query produces when it does match rows. Consumers therefore never have /// to reconstruct a schema the segment already knows. private void assertZeroMatchEmitsOneSchemaBlock(IndexSegment segment, @Language("sql") String zeroMatchQuery, @Language("sql") String matchingQuery, int numSortedExpressions, boolean nullHandling) { @@ -592,11 +592,11 @@ public class StreamingSelectionOrderByOperatorTest { assertNull(operator.nextBlock(), "Exactly one block may precede exhaustion"); } - /// Drives the operator twice over {@code segment}: once against the real project operator, once against one that + /// Drives the operator twice over `segment`: once against the real project operator, once against one that /// splices zero-document blocks into its output, and asserts the emitted rows are identical. /// - /// <p>No doc-id-set operator produces an empty block today, so the case has to be injected. The project operator is - /// built with a deliberately small {@code maxDocsPerCall} so the scan spans several blocks. + /// No doc-id-set operator produces an empty block today, so the case has to be injected. The project operator is + /// built with a deliberately small `maxDocsPerCall` so the scan spans several blocks. /// /// @param injectBeforeRealBlock 1-based index of the real block to splice empties in front of; 0 for none /// @param numEmptyBlocks how many consecutive empty blocks to splice in at that point @@ -625,7 +625,7 @@ public class StreamingSelectionOrderByOperatorTest { } } - /// Builds a {@link StreamingSelectionOrderByOperator} directly (rather than through {@link SelectionPlanNode}) so + /// Builds a [StreamingSelectionOrderByOperator] directly (rather than through [SelectionPlanNode]) so /// the project operator can be wrapped. Passing 0 for every injection parameter yields the undecorated operator. private StreamingSelectionOrderByOperator buildOperatorWithInjectableProject(IndexSegment segment, @Language("sql") String query, int numSortedExpressions, boolean nullHandling, int injectBeforeRealBlock, @@ -665,10 +665,10 @@ public class StreamingSelectionOrderByOperatorTest { } /// Forwards everything to a real project operator, but splices zero-document blocks into its output: - /// {@code numEmptyBlocks} of them just before the {@code injectBeforeRealBlock}-th real block, and - /// {@code numTrailingEmptyBlocks} after the last real block but before exhaustion. + /// `numEmptyBlocks` of them just before the `injectBeforeRealBlock`-th real block, and + /// `numTrailingEmptyBlocks` after the last real block but before exhaustion. /// - /// <p>Every real block is still delivered, merely deferred, so the decorator provably drops nothing of its own - + /// Every real block is still delivered, merely deferred, so the decorator provably drops nothing of its own - /// any row loss observed by a test is the operator under test truncating its scan. private static class EmptyBlockInjectingProjectOperator extends BaseProjectOperator<ValueBlock> { private final BaseProjectOperator<?> _delegate; @@ -791,13 +791,13 @@ public class StreamingSelectionOrderByOperatorTest { } } - /// Runs {@code query} twice over {@code segment} - once with the streaming hint on, once off - and asserts the + /// Runs `query` twice over `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}. + /// cases that span multiple runs pass `>= 2` to prove the output is genuinely streamed; cases whose + /// result fits in a single trimmed block pass `1`. private StreamingSelectionOrderByOperator assertParity(IndexSegment segment, @Language("sql") String query, boolean nullHandling, int expectedMinBlocks) { // Streaming path. @@ -842,7 +842,7 @@ public class StreamingSelectionOrderByOperatorTest { return (StreamingSelectionOrderByOperator) streamingOperator; } - /// Drains the streaming operator for {@code query} and returns all rows concatenated in emission order. + /// Drains the streaming operator for `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); diff --git a/pinot-core/src/test/java/org/apache/pinot/core/plan/maker/InstancePlanMakerImplV2Test.java b/pinot-core/src/test/java/org/apache/pinot/core/plan/maker/InstancePlanMakerImplV2Test.java index 64c6126bc43..45270526ea1 100644 --- a/pinot-core/src/test/java/org/apache/pinot/core/plan/maker/InstancePlanMakerImplV2Test.java +++ b/pinot-core/src/test/java/org/apache/pinot/core/plan/maker/InstancePlanMakerImplV2Test.java @@ -29,11 +29,12 @@ import org.testng.annotations.Test; import static org.testng.Assert.assertEquals; -/// Tests for execution-thread resolution in [InstancePlanMakerImplV2], covering the interplay -/// between `max.execution.threads`, `default.execution.threads`, and per-query overrides. +/// Tests for server-config versus query-option resolution in [InstancePlanMakerImplV2]: execution threads +/// (`max.execution.threads`, `default.execution.threads`) and the streaming selection ORDER BY merge defaults. public class InstancePlanMakerImplV2Test { private static final String BASE_QUERY = "SELECT * FROM testTable"; + private static final String SELECTION_ORDER_BY_QUERY = "SELECT col FROM testTable ORDER BY col LIMIT 10"; private static QueryContext buildQueryContext(String maxExecutionThreadsOption) { QueryContext queryContext = QueryContextConverterUtils.getQueryContext(BASE_QUERY); @@ -99,4 +100,66 @@ public class InstancePlanMakerImplV2Test { config.setProperty(Server.DEFAULT_EXECUTION_THREADS, 8); planMaker.init(config); } + + @Test + public void testSortedSelectionMergeDefaultsWithoutServerConfig() { + InstancePlanMakerImplV2 planMaker = new InstancePlanMakerImplV2(); + planMaker.init(new PinotConfiguration()); + + QueryContext queryContext = QueryContextConverterUtils.getQueryContext(SELECTION_ORDER_BY_QUERY); + planMaker.applyQueryOptions(queryContext); + + assertEquals(queryContext.getSortedSelectionMergeAutoMinSortedRatio(), + Server.DEFAULT_SORTED_SELECTION_MERGE_AUTO_MIN_SORTED_RATIO); + assertEquals(queryContext.getSortedSelectionMergeBlockSize(), Server.DEFAULT_SORTED_SELECTION_MERGE_BLOCK_SIZE); + } + + @Test + public void testSortedSelectionMergeServerConfigAppliesWithoutQueryOption() { + InstancePlanMakerImplV2 planMaker = new InstancePlanMakerImplV2(); + planMaker.init(sortedSelectionMergeConfig(0.6, 2_500)); + + QueryContext queryContext = QueryContextConverterUtils.getQueryContext(SELECTION_ORDER_BY_QUERY); + planMaker.applyQueryOptions(queryContext); + + assertEquals(queryContext.getSortedSelectionMergeAutoMinSortedRatio(), 0.6); + assertEquals(queryContext.getSortedSelectionMergeBlockSize(), 2_500); + } + + @Test + public void testSortedSelectionMergeQueryOptionOverridesServerConfig() { + InstancePlanMakerImplV2 planMaker = new InstancePlanMakerImplV2(); + planMaker.init(sortedSelectionMergeConfig(0.6, 2_500)); + + QueryContext queryContext = QueryContextConverterUtils.getQueryContext(SELECTION_ORDER_BY_QUERY); + Map<String, String> queryOptions = queryContext.getQueryOptions(); + queryOptions.put(QueryOptionKey.SORTED_SELECTION_MERGE_AUTO_MIN_SORTED_RATIO, "0.9"); + queryOptions.put(QueryOptionKey.SORTED_SELECTION_MERGE_BLOCK_SIZE, "500"); + planMaker.applyQueryOptions(queryContext); + + assertEquals(queryContext.getSortedSelectionMergeAutoMinSortedRatio(), 0.9); + assertEquals(queryContext.getSortedSelectionMergeBlockSize(), 500); + } + + @Test(expectedExceptions = IllegalStateException.class) + public void testInitRejectsSortedSelectionMergeRatioAboveOne() { + new InstancePlanMakerImplV2().init(sortedSelectionMergeConfig(1.5, 2_500)); + } + + @Test(expectedExceptions = IllegalStateException.class) + public void testInitRejectsNegativeSortedSelectionMergeRatio() { + new InstancePlanMakerImplV2().init(sortedSelectionMergeConfig(-0.1, 2_500)); + } + + @Test(expectedExceptions = IllegalStateException.class) + public void testInitRejectsNonPositiveSortedSelectionMergeBlockSize() { + new InstancePlanMakerImplV2().init(sortedSelectionMergeConfig(0.6, 0)); + } + + private static PinotConfiguration sortedSelectionMergeConfig(double autoMinSortedRatio, int blockSize) { + PinotConfiguration config = new PinotConfiguration(); + config.setProperty(Server.SORTED_SELECTION_MERGE_AUTO_MIN_SORTED_RATIO, autoMinSortedRatio); + config.setProperty(Server.SORTED_SELECTION_MERGE_BLOCK_SIZE, blockSize); + return config; + } } 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 4a8d90207b1..efbbbe15f1a 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 @@ -631,10 +631,6 @@ public class CommonConstants { "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"; @@ -1650,6 +1646,21 @@ public class CommonConstants { /// among many sorted ones. This default is chosen to tolerate that while still rejecting a mostly-unsorted table. /// It is a starting point rather than a measured optimum. public static final double DEFAULT_SORTED_SELECTION_MERGE_AUTO_MIN_SORTED_RATIO = 0.8; + /// Server-wide override of [#DEFAULT_SORTED_SELECTION_MERGE_AUTO_MIN_SORTED_RATIO], applied to queries that do not + /// set the query option. + public static final String SORTED_SELECTION_MERGE_AUTO_MIN_SORTED_RATIO = + "sorted.selection.merge.auto.min.sorted.ratio"; + public static final String CONFIG_OF_QUERY_EXECUTOR_SORTED_SELECTION_MERGE_AUTO_MIN_SORTED_RATIO = + QUERY_EXECUTOR_CONFIG_PREFIX + "." + SORTED_SELECTION_MERGE_AUTO_MIN_SORTED_RATIO; + + /// Default output block size (rows) for the streaming selection ORDER BY combine + /// ([Broker.Request.QueryOptionKey#SORTED_SELECTION_MERGE_BLOCK_SIZE]). + public static final int DEFAULT_SORTED_SELECTION_MERGE_BLOCK_SIZE = 10_000; + /// Server-wide override of [#DEFAULT_SORTED_SELECTION_MERGE_BLOCK_SIZE], applied to queries that do not set the + /// query option. + public static final String SORTED_SELECTION_MERGE_BLOCK_SIZE = "sorted.selection.merge.block.size"; + public static final String CONFIG_OF_QUERY_EXECUTOR_SORTED_SELECTION_MERGE_BLOCK_SIZE = + QUERY_EXECUTOR_CONFIG_PREFIX + "." + SORTED_SELECTION_MERGE_BLOCK_SIZE; public static final String CONFIG_OF_MSE_MIN_GROUP_TRIM_SIZE = MSE_CONFIG_PREFIX + ".min.group.trim.size"; // Match the value of GroupByUtils.DEFAULT_MIN_NUM_GROUPS public static final int DEFAULT_MSE_MIN_GROUP_TRIM_SIZE = 5000; --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
