Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -105,7 +105,7 @@ public void testOrderByKeysIsPushedToFinalAggregationStageWhenGroupTrimIsEnabled
final String trimEnabledPlan = "Execution Plan\n"
+ "LogicalSort(sort0=[$0], sort1=[$1], dir0=[ASC], dir1=[ASC], offset=[0], fetch=[3])\n"
+ " PinotLogicalSortExchange(distribution=[hash], collation=[[0, 1]], isSortOnSender=[false], "
+ "isSortOnReceiver=[true])\n"
+ "isSortOnReceiver=[false])\n"
+ " LogicalSort(sort0=[$0], sort1=[$1], dir0=[ASC], dir1=[ASC], fetch=[3])\n"
// 'collations' below is the important bit
+ " PinotLogicalAggregate(group=[{0, 1}], agg#0=[COUNT($2)], aggType=[FINAL], collations=[[0, "
Expand Down Expand Up @@ -164,7 +164,7 @@ public void testOrderByKeysIsNotPushedToFinalAggregationStageWhenGroupTrimIsDisa
"Execution Plan\n"
+ "LogicalSort(sort0=[$0], sort1=[$1], dir0=[ASC], dir1=[ASC], offset=[0], fetch=[3])\n"
+ " PinotLogicalSortExchange(distribution=[hash], collation=[[0, 1]], isSortOnSender=[false], "
+ "isSortOnReceiver=[true])\n"
+ "isSortOnReceiver=[false])\n"
+ " LogicalSort(sort0=[$0], sort1=[$1], dir0=[ASC], dir1=[ASC], fetch=[3])\n"
// lack of 'collations' below is the important bit
+ " PinotLogicalAggregate(group=[{0, 1}], agg#0=[COUNT($2)], aggType=[FINAL])\n"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3657,7 +3657,7 @@ public void testExplainPlanQueryV2()
assertEquals(response1Json.get("rows").get(0).get(1).asText(), "Execution Plan\n"
+ "LogicalSort(sort0=[$0], dir0=[ASC])\n"
+ " PinotLogicalSortExchange(distribution=[hash], collation=[[0]], isSortOnSender=[false], "
+ "isSortOnReceiver=[true])\n"
+ "isSortOnReceiver=[false])\n"
+ " LogicalProject(count=[$1], name=[$0])\n"
+ " PinotLogicalAggregate(group=[{0}], agg#0=[COUNT($1)], aggType=[FINAL])\n"
+ " PinotLogicalExchange(distribution=[hash[0]])\n"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -127,7 +127,7 @@ public void testOrderByKeysIsNotPushedToFinalAggregationWhenGroupTrimHintIsDisab
String trimDisabledPlan = "Execution Plan\n"
+ "LogicalSort(sort0=[$0], sort1=[$1], dir0=[ASC], dir1=[DESC], offset=[0], fetch=[1])\n"
+ " PinotLogicalSortExchange(distribution=[hash], collation=[[0, 1 DESC]], isSortOnSender=[false], "
+ "isSortOnReceiver=[true])\n"
+ "isSortOnReceiver=[false])\n"
+ " LogicalSort(sort0=[$0], sort1=[$1], dir0=[ASC], dir1=[DESC], fetch=[1])\n"
+ " PinotLogicalAggregate(group=[{0, 1}], agg#0=[COUNT($2)], aggType=[FINAL])\n"
+ " PinotLogicalExchange(distribution=[hash[0, 1]])\n"
Expand Down Expand Up @@ -188,7 +188,7 @@ public void testOrderByKeysIsPushedToFinalAggregationStageWithoutGroupTrimSize()
"Execution Plan\n"
+ "LogicalSort(sort0=[$0], sort1=[$1], dir0=[ASC], dir1=[DESC], offset=[0], fetch=[1])\n"
+ " PinotLogicalSortExchange(distribution=[hash], collation=[[0, 1 DESC]], isSortOnSender=[false], "
+ "isSortOnReceiver=[true])\n"
+ "isSortOnReceiver=[false])\n"
+ " LogicalSort(sort0=[$0], sort1=[$1], dir0=[ASC], dir1=[DESC], fetch=[1])\n"
+ " PinotLogicalAggregate(group=[{0, 1}], agg#0=[COUNT($2)], aggType=[FINAL], collations=[[0, 1 "
+ "DESC]], limit=[1])\n"
Expand Down Expand Up @@ -224,7 +224,7 @@ public void testOrderByKeysIsPushedToFinalAggregationStageWithGroupTrimSize()
"Execution Plan\n"
+ "LogicalSort(sort0=[$0], sort1=[$1], dir0=[ASC], dir1=[DESC], offset=[0], fetch=[1])\n"
+ " PinotLogicalSortExchange(distribution=[hash], collation=[[0, 1 DESC]], isSortOnSender=[false], "
+ "isSortOnReceiver=[true])\n"
+ "isSortOnReceiver=[false])\n"
+ " LogicalSort(sort0=[$0], sort1=[$1], dir0=[ASC], dir1=[DESC], fetch=[1])\n"
+ " PinotLogicalAggregate(group=[{0, 1}], agg#0=[COUNT($2)], aggType=[FINAL], collations=[[0, 1 "
+ "DESC]], limit=[1])\n"
Expand Down Expand Up @@ -258,7 +258,7 @@ public void testOrderByKeysIsPushedToFinalAggregationStage()
"Execution Plan\n"
+ "LogicalSort(sort0=[$0], sort1=[$1], dir0=[ASC], dir1=[ASC], offset=[0], fetch=[3])\n"
+ " PinotLogicalSortExchange(distribution=[hash], collation=[[0, 1]], isSortOnSender=[false], "
+ "isSortOnReceiver=[true])\n"
+ "isSortOnReceiver=[false])\n"
+ " LogicalSort(sort0=[$0], sort1=[$1], dir0=[ASC], dir1=[ASC], fetch=[3])\n"
+ " PinotLogicalAggregate(group=[{0, 1}], agg#0=[COUNT($2)], aggType=[FINAL], collations=[[0, "
+ "1]], limit=[3])\n"
Expand Down Expand Up @@ -293,7 +293,7 @@ public void testHavingOnKeysAndOrderByKeysIsPushedToFinalAggregationStage()
"Execution Plan\n"
+ "LogicalSort(sort0=[$0], sort1=[$1], dir0=[ASC], dir1=[ASC], offset=[0], fetch=[3])\n"
+ " PinotLogicalSortExchange(distribution=[hash], collation=[[0, 1]], isSortOnSender=[false], "
+ "isSortOnReceiver=[true])\n"
+ "isSortOnReceiver=[false])\n"
+ " LogicalSort(sort0=[$0], sort1=[$1], dir0=[ASC], dir1=[ASC], fetch=[3])\n"
+ " PinotLogicalAggregate(group=[{0, 1}], agg#0=[COUNT($2)], aggType=[FINAL], collations=[[0, "
+ "1]], limit=[3])\n"
Expand Down Expand Up @@ -328,7 +328,7 @@ public void testGroupByKeysWithOffsetIsPushedToFinalAggregationStage()
"Execution Plan\n"
+ "LogicalSort(sort0=[$0], sort1=[$1], dir0=[ASC], dir1=[ASC], offset=[1], fetch=[3])\n"
+ " PinotLogicalSortExchange(distribution=[hash], collation=[[0, 1]], isSortOnSender=[false], "
+ "isSortOnReceiver=[true])\n"
+ "isSortOnReceiver=[false])\n"
+ " LogicalSort(sort0=[$0], sort1=[$1], dir0=[ASC], dir1=[ASC], fetch=[4])\n"
+ " PinotLogicalAggregate(group=[{0, 1}], agg#0=[COUNT($2)], aggType=[FINAL], collations=[[0, "
+ "1]], limit=[4])\n"
Expand Down Expand Up @@ -470,7 +470,7 @@ public void testOrderByByKeysAndValuesIsPushedToFinalAggregationStage()
+ "LogicalSort(sort0=[$0], sort1=[$1], sort2=[$2], dir0=[DESC], dir1=[DESC], dir2=[DESC], offset=[0],"
+ " fetch=[3])\n"
+ " PinotLogicalSortExchange(distribution=[hash], collation=[[0 DESC, 1 DESC, 2 DESC]], "
+ "isSortOnSender=[false], isSortOnReceiver=[true])\n"
+ "isSortOnSender=[false], isSortOnReceiver=[false])\n"
+ " LogicalSort(sort0=[$0], sort1=[$1], sort2=[$2], dir0=[DESC], dir1=[DESC], dir2=[DESC], "
+ "fetch=[3])\n"
+ " PinotLogicalAggregate(group=[{0, 1}], agg#0=[COUNT($2)], aggType=[FINAL], collations=[[0 "
Expand Down Expand Up @@ -511,7 +511,7 @@ public void testOrderByKeyValueExpressionIsNotPushedToFinalAggregateStage()
"Execution Plan\n"
+ "LogicalSort(sort0=[$3], dir0=[DESC], offset=[0], fetch=[3])\n"
+ " PinotLogicalSortExchange(distribution=[hash], collation=[[3 DESC]], isSortOnSender=[false], "
+ "isSortOnReceiver=[true])\n"
+ "isSortOnReceiver=[false])\n"
+ " LogicalSort(sort0=[$3], dir0=[DESC], fetch=[3])\n"
+ " LogicalProject(i=[$0], j=[$1], cnt=[$2], EXPR$3=[*(*($0, $1), $2)])\n"
+ " PinotLogicalAggregate(group=[{0, 1}], agg#0=[COUNT($2)], aggType=[FINAL])\n"
Expand Down Expand Up @@ -549,7 +549,7 @@ public void testForGroupByOverJoinOrderByKeyIsPushedToAggregationLeafStage()
"Execution Plan\n"
+ "LogicalSort(sort0=[$0], sort1=[$1], dir0=[ASC], dir1=[ASC], offset=[0], fetch=[5])\n"
+ " PinotLogicalSortExchange(distribution=[hash], collation=[[0, 1]], isSortOnSender=[false], "
+ "isSortOnReceiver=[true])\n"
+ "isSortOnReceiver=[false])\n"
+ " LogicalSort(sort0=[$0], sort1=[$1], dir0=[ASC], dir1=[ASC], fetch=[5])\n"
+ " PinotLogicalAggregate(group=[{0, 1}], agg#0=[COUNT($2)], aggType=[FINAL], collations=[[0, "
+ "1]], limit=[5])\n"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -143,7 +143,7 @@ public void testMSQEOrderByOnDependingOnAggregateResultIsNotPushedDown()
"Execution Plan\n"
+ "LogicalSort(sort0=[$3], dir0=[DESC], offset=[0], fetch=[5])\n"
+ " PinotLogicalSortExchange(distribution=[hash], collation=[[3 DESC]], isSortOnSender=[false], "
+ "isSortOnReceiver=[true])\n"
+ "isSortOnReceiver=[false])\n"
+ " LogicalSort(sort0=[$3], dir0=[DESC], fetch=[5])\n" // <-- actual sort & limit
+ " LogicalProject(i=[$0], j=[$1], EXPR$2=[$2], EXPR$3=[*($0, $1)])\n"
// <-- order by value is computed here, so trimming in upstream stages is not possible
Expand Down Expand Up @@ -186,7 +186,7 @@ public void testMSQEGroupsTrimmedAtSegmentLevelWithOrderByOnSomeGroupByKeysIsNot
"Execution Plan\n"
+ "LogicalSort(sort0=[$1], dir0=[DESC], offset=[0], fetch=[5])\n"
+ " PinotLogicalSortExchange(distribution=[hash], collation=[[1 DESC]], isSortOnSender=[false], "
+ "isSortOnReceiver=[true])\n"
+ "isSortOnReceiver=[false])\n"
+ " LogicalSort(sort0=[$1], dir0=[DESC], fetch=[5])\n"
+ " PinotLogicalAggregate(group=[{0, 1}], agg#0=[COUNT($2)], aggType=[FINAL], collations=[[1 DESC]],"
+ " limit=[5])\n"
Expand Down Expand Up @@ -237,7 +237,7 @@ public void testMSQEGroupsTrimmedAtSegmentLevelWithOrderByOnAggregateIsNotSafe()
"Execution Plan\n"
+ "LogicalSort(sort0=[$2], dir0=[DESC], offset=[0], fetch=[5])\n"
+ " PinotLogicalSortExchange(distribution=[hash], collation=[[2 DESC]], isSortOnSender=[false], "
+ "isSortOnReceiver=[true])\n"
+ "isSortOnReceiver=[false])\n"
+ " LogicalSort(sort0=[$2], dir0=[DESC], fetch=[5])\n"
+ " PinotLogicalAggregate(group=[{0, 1}], agg#0=[COUNT($2)], aggType=[FINAL], collations=[[2 DESC]],"
+ " limit=[5])\n"
Expand Down Expand Up @@ -281,7 +281,7 @@ public void testMSQEGroupsTrimmedAtInterSegmentLevelWithOrderByOnSomeGroupByKeys
"Execution Plan\n"
+ "LogicalSort(sort0=[$1], dir0=[DESC], offset=[0], fetch=[5])\n"
+ " PinotLogicalSortExchange(distribution=[hash], collation=[[1 DESC]], isSortOnSender=[false], "
+ "isSortOnReceiver=[true])\n"
+ "isSortOnReceiver=[false])\n"
+ " LogicalSort(sort0=[$1], dir0=[DESC], fetch=[5])\n"
+ " PinotLogicalAggregate(group=[{0, 1}], agg#0=[COUNT($2)], aggType=[FINAL], collations=[[1 DESC]],"
+ " limit=[5])\n"
Expand Down Expand Up @@ -326,7 +326,7 @@ public void testMSQEGroupsTrimmedAtIntermediateLevelWithOrderByOnSomeGroupByKeys
"Execution Plan\n"
+ "LogicalSort(sort0=[$1], dir0=[DESC], offset=[0], fetch=[5])\n"
+ " PinotLogicalSortExchange(distribution=[hash], collation=[[1 DESC]], isSortOnSender=[false], "
+ "isSortOnReceiver=[true])\n"
+ "isSortOnReceiver=[false])\n"
+ " LogicalSort(sort0=[$1], dir0=[DESC], fetch=[5])\n"
+ " PinotLogicalAggregate(group=[{0, 1}], agg#0=[COUNT($2)], aggType=[FINAL], collations=[[1 DESC]],"
+ " limit=[5])\n" // receives 50-row-big blocks, trimming kicks in only if limit is lower
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,182 @@
/**
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/
package org.apache.pinot.perf;

import java.util.ArrayList;
import java.util.Arrays;
import java.util.Comparator;
import java.util.List;
import java.util.PriorityQueue;
import java.util.Random;
import java.util.concurrent.TimeUnit;
import org.apache.calcite.rel.RelFieldCollation;
import org.apache.pinot.core.query.selection.SelectionOperatorUtils;
import org.apache.pinot.query.runtime.operator.utils.SortUtils;
import org.openjdk.jmh.annotations.Benchmark;
import org.openjdk.jmh.annotations.BenchmarkMode;
import org.openjdk.jmh.annotations.Fork;
import org.openjdk.jmh.annotations.Measurement;
import org.openjdk.jmh.annotations.Mode;
import org.openjdk.jmh.annotations.OutputTimeUnit;
import org.openjdk.jmh.annotations.Param;
import org.openjdk.jmh.annotations.Scope;
import org.openjdk.jmh.annotations.Setup;
import org.openjdk.jmh.annotations.State;
import org.openjdk.jmh.annotations.Warmup;
import org.openjdk.jmh.infra.Blackhole;
import org.openjdk.jmh.runner.Runner;
import org.openjdk.jmh.runner.RunnerException;
import org.openjdk.jmh.runner.options.OptionsBuilder;


/// Compares the two ways the multi-stage engine can produce a fully sorted result when nothing bounds it - i.e. an
/// ORDER BY with no LIMIT.
///
/// - [#unboundedHeap] reproduces the SortOperator path before the split: an unbounded PriorityQueue fed one row at
/// a time through [SelectionOperatorUtils#addToPriorityQueue], then drained by polling every row. Both halves
/// pay a sift.
/// - [#accumulateAndSort] reproduces FullSortOperator: append every row to an ArrayList, then sort once.
/// - [#accumulateAndSortNotPreSized] is the same, from a default-capacity list, to isolate what pre-sizing buys.
///
/// Results at the time of writing (3 forks x 10 iterations, 3-column rows), for orientation only - re-run rather
/// than trust these:
///
/// | input | rows | sort | heap | speedup |
/// |----------------|------|--------|--------|---------|
/// | random | 10K | 0.70ms | 1.09ms | 1.6x |
/// | random | 1M | 302ms | 590ms | 2.0x |
/// | 64 sorted runs | 10K | 0.23ms | 0.91ms | 4.0x |
/// | 64 sorted runs | 1M | 82ms | 369ms | 4.5x |
///
/// The sort allocates 1.2x (1M rows) to 1.7x (10K rows) what the heap does: the merge buffer is the difference.
///
/// Two input distributions are measured, because the shape of the input is what decides whether the comparison is
/// close. RANDOM is the neutral case. SORTED_RUNS is what a receive stage actually sees once the senders sort - the
/// concatenation of one sorted run per sender - and TimSort detects such runs while a heap cannot.
@BenchmarkMode(Mode.AverageTime)
@OutputTimeUnit(TimeUnit.MILLISECONDS)
@Fork(1)
@Warmup(iterations = 3, time = 1)
@Measurement(iterations = 5, time = 1)
@State(Scope.Benchmark)
public class BenchmarkMseSortImplementations {
private static final int NUM_COLUMNS = 3;
private static final int NUM_SENDERS = 64;
/// Matches SelectionOperatorUtils.MAX_ROW_HOLDER_INITIAL_CAPACITY, which the old operator used.
private static final int HOLDER_CAPACITY = 10_000;
/// Rows per arriving block, so the accumulation pattern matches what an operator actually sees.
private static final int ROWS_PER_BLOCK = 2_048;

@Param({"10000", "1000000"})
private int _numRows;

@Param({"RANDOM", "SORTED_RUNS"})
private String _distribution;

private Object[][] _rows;
private Comparator<Object[]> _forward;
private Comparator<Object[]> _reversed;

@Setup
public void setUp() {
List<RelFieldCollation> collations = List.of(
new RelFieldCollation(0, RelFieldCollation.Direction.ASCENDING, RelFieldCollation.NullDirection.LAST));
_forward = new SortUtils.SortComparator(collations, false);
// The old path used the inverted comparator so the heap head is the row to evict.
_reversed = new SortUtils.SortComparator(collations, true);

Random random = new Random(42);
_rows = new Object[_numRows][];
if ("RANDOM".equals(_distribution)) {
for (int i = 0; i < _numRows; i++) {
_rows[i] = row(random.nextInt(), i);
}
} else {
// One sorted run per sender, concatenated in mailbox order: what the receiver sees when senders sort.
int perSender = (_numRows + NUM_SENDERS - 1) / NUM_SENDERS;
int written = 0;
for (int sender = 0; sender < NUM_SENDERS && written < _numRows; sender++) {
int[] keys = new int[Math.min(perSender, _numRows - written)];
for (int i = 0; i < keys.length; i++) {
keys[i] = random.nextInt();
}
Arrays.sort(keys);
for (int key : keys) {
_rows[written] = row(key, written);
written++;
}
}
}
}

private static Object[] row(int key, int seq) {
Object[] row = new Object[NUM_COLUMNS];
row[0] = key;
row[1] = (long) seq;
row[2] = seq;
return row;
}

/// The SortOperator path before the split, with `numRowsToKeep == Integer.MAX_VALUE` so nothing is ever evicted.
@Benchmark
public void unboundedHeap(Blackhole bh) {
PriorityQueue<Object[]> queue = new PriorityQueue<>(HOLDER_CAPACITY, _reversed);
for (Object[] row : _rows) {
SelectionOperatorUtils.addToPriorityQueue(row, queue, Integer.MAX_VALUE);
}
int resultSize = queue.size();
Object[][] result = new Object[resultSize][];
for (int i = resultSize - 1; i >= 0; i--) {
result[i] = queue.poll();
}
bh.consume(result);
}

/// FullSortOperator: accumulate into an ArrayList, sort once.
///
/// Rows arrive one block at a time, so this appends in block-sized chunks rather than in one addAll - that is what
/// decides whether the backing array grows repeatedly, and therefore whether pre-sizing the list is worth anything.
@Benchmark
public void accumulateAndSort(Blackhole bh) {
ArrayList<Object[]> rows = new ArrayList<>(HOLDER_CAPACITY);
for (int start = 0; start < _rows.length; start += ROWS_PER_BLOCK) {
int end = Math.min(start + ROWS_PER_BLOCK, _rows.length);
rows.addAll(Arrays.asList(_rows).subList(start, end));
}
rows.sort(_forward);
bh.consume(rows);
}

/// Same, but from a default-capacity list, to isolate what pre-sizing buys.
@Benchmark
public void accumulateAndSortNotPreSized(Blackhole bh) {
ArrayList<Object[]> rows = new ArrayList<>();
for (int start = 0; start < _rows.length; start += ROWS_PER_BLOCK) {
int end = Math.min(start + ROWS_PER_BLOCK, _rows.length);
rows.addAll(Arrays.asList(_rows).subList(start, end));
}
rows.sort(_forward);
bh.consume(rows);
}

public static void main(String[] args)
throws RunnerException {
new Runner(new OptionsBuilder().include(BenchmarkMseSortImplementations.class.getSimpleName()).build()).run();
}
}
Loading
Loading