This is an automated email from the ASF dual-hosted git repository.
yashmayya pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pinot.git
The following commit(s) were added to refs/heads/master by this push:
new fcbe3e4a73b Seal large IN lists from Calcite during multi-stage
planning (#19674)
fcbe3e4a73b is described below
commit fcbe3e4a73b7f05b73b836dd96c49252d3dd1d59
Author: Yash Mayya <[email protected]>
AuthorDate: Wed Sep 30 07:57:30 2026 -0700
Seal large IN lists from Calcite during multi-stage planning (#19674)
---
.../MultiStageBrokerRequestHandler.java | 4 +
.../common/utils/config/QueryOptionsUtils.java | 8 +
.../common/utils/config/QueryOptionsUtilsTest.java | 25 +
.../pinot/calcite/rel/rules/PinotRuleUtils.java | 25 +
.../calcite/rex/PinotSealedSearchOperator.java | 117 ++++
.../org/apache/pinot/calcite/rex/SearchSealer.java | 401 ++++++++++++++
.../pinot/calcite/sql/fun/PinotInListOperator.java | 93 ++++
.../calcite/sql2rel/PinotConvertletTable.java | 140 ++++-
.../org/apache/pinot/query/QueryEnvironment.java | 17 +-
.../apache/pinot/query/context/PlannerContext.java | 15 +
.../query/planner/logical/RexExpressionUtils.java | 103 +++-
.../apache/pinot/calcite/rex/SearchSealerTest.java | 315 +++++++++++
.../pinot/query/SealedInListPlanningTest.java | 555 +++++++++++++++++++
.../planner/logical/RexExpressionUtilsTest.java | 149 +++++
.../src/test/resources/queries/LargeInLists.json | 602 +++++++++++++++++++++
.../apache/pinot/spi/utils/CommonConstants.java | 23 +
16 files changed, 2580 insertions(+), 12 deletions(-)
diff --git
a/pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/MultiStageBrokerRequestHandler.java
b/pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/MultiStageBrokerRequestHandler.java
index 73d348dc437..b3e8f6f18a8 100644
---
a/pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/MultiStageBrokerRequestHandler.java
+++
b/pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/MultiStageBrokerRequestHandler.java
@@ -581,6 +581,9 @@ public class MultiStageBrokerRequestHandler extends
BaseBrokerRequestHandler {
int sortExchangeCopyThreshold = _config.getProperty(
CommonConstants.Broker.CONFIG_OF_SORT_EXCHANGE_COPY_THRESHOLD,
CommonConstants.Broker.DEFAULT_SORT_EXCHANGE_COPY_THRESHOLD);
+ int sealedInListThreshold = _config.getProperty(
+ CommonConstants.Broker.CONFIG_OF_SEALED_IN_LIST_THRESHOLD,
+ CommonConstants.Broker.DEFAULT_SEALED_IN_LIST_THRESHOLD);
boolean defaultUnnestColumnPruning = _config.getProperty(
CommonConstants.Broker.CONFIG_OF_UNNEST_COLUMN_PRUNING,
CommonConstants.Broker.DEFAULT_UNNEST_COLUMN_PRUNING);
@@ -619,6 +622,7 @@ public class MultiStageBrokerRequestHandler extends
BaseBrokerRequestHandler {
.defaultHashFunction(defaultHashFunction)
.defaultDisabledPlannerRules(_defaultDisabledPlannerRules)
.defaultSortExchangeCopyLimit(sortExchangeCopyThreshold)
+ .defaultSealedInListThreshold(sealedInListThreshold)
.build();
}
diff --git
a/pinot-common/src/main/java/org/apache/pinot/common/utils/config/QueryOptionsUtils.java
b/pinot-common/src/main/java/org/apache/pinot/common/utils/config/QueryOptionsUtils.java
index 76519cc1cc4..3b0da6329eb 100644
---
a/pinot-common/src/main/java/org/apache/pinot/common/utils/config/QueryOptionsUtils.java
+++
b/pinot-common/src/main/java/org/apache/pinot/common/utils/config/QueryOptionsUtils.java
@@ -1035,4 +1035,12 @@ public class QueryOptionsUtils {
}
return i;
}
+
+ /// Returns the [QueryOptionKey#SEALED_IN_LIST_THRESHOLD] option, or
`defaultValue` when the option is not set.
+ public static int getSealedInListThreshold(Map<String, String> options, int
defaultValue) {
+ String threshold = options.get(QueryOptionKey.SEALED_IN_LIST_THRESHOLD);
+ Integer value =
+ uncheckedParseInt(QueryOptionKey.SEALED_IN_LIST_THRESHOLD, threshold
!= null ? threshold.trim() : null);
+ return value != null ? value : defaultValue;
+ }
}
diff --git
a/pinot-common/src/test/java/org/apache/pinot/common/utils/config/QueryOptionsUtilsTest.java
b/pinot-common/src/test/java/org/apache/pinot/common/utils/config/QueryOptionsUtilsTest.java
index 29dad29d5b7..73d248c5c09 100644
---
a/pinot-common/src/test/java/org/apache/pinot/common/utils/config/QueryOptionsUtilsTest.java
+++
b/pinot-common/src/test/java/org/apache/pinot/common/utils/config/QueryOptionsUtilsTest.java
@@ -33,6 +33,7 @@ import static org.testng.Assert.assertEquals;
import static org.testng.Assert.assertFalse;
import static org.testng.Assert.assertNull;
import static org.testng.Assert.assertTrue;
+import static org.testng.Assert.expectThrows;
import static org.testng.Assert.fail;
@@ -410,4 +411,28 @@ public class QueryOptionsUtilsTest {
queryOptions.put(LITE_MODE_IMPLICIT_LEAF_STAGE_LIMIT, "0");
assertEquals(QueryOptionsUtils.getLiteModeImplicitLeafStageLimit(queryOptions),
Integer.valueOf(0));
}
+
+ @Test
+ public void testGetSealedInListThreshold() {
+ Map<String, String> queryOptions = new HashMap<>();
+
+ // Absent → default
+ assertEquals(QueryOptionsUtils.getSealedInListThreshold(queryOptions, 20),
20);
+
+ // Present → parsed value
+ queryOptions.put(SEALED_IN_LIST_THRESHOLD, " 100 ");
+ assertEquals(QueryOptionsUtils.getSealedInListThreshold(queryOptions, 20),
100);
+
+ // Zero or negative → returned as is, which turns sealing off
+ queryOptions.put(SEALED_IN_LIST_THRESHOLD, "0");
+ assertEquals(QueryOptionsUtils.getSealedInListThreshold(queryOptions, 20),
0);
+ queryOptions.put(SEALED_IN_LIST_THRESHOLD, "-1");
+ assertEquals(QueryOptionsUtils.getSealedInListThreshold(queryOptions, 20),
-1);
+
+ // Not an integer → error that names the option
+ queryOptions.put(SEALED_IN_LIST_THRESHOLD, "twenty");
+ IllegalArgumentException e = expectThrows(IllegalArgumentException.class,
+ () -> QueryOptionsUtils.getSealedInListThreshold(queryOptions, 20));
+ assertTrue(e.getMessage().contains(SEALED_IN_LIST_THRESHOLD),
e.getMessage());
+ }
}
diff --git
a/pinot-query-planner/src/main/java/org/apache/pinot/calcite/rel/rules/PinotRuleUtils.java
b/pinot-query-planner/src/main/java/org/apache/pinot/calcite/rel/rules/PinotRuleUtils.java
index b7ad7a4edb2..bb7f98a24c4 100644
---
a/pinot-query-planner/src/main/java/org/apache/pinot/calcite/rel/rules/PinotRuleUtils.java
+++
b/pinot-query-planner/src/main/java/org/apache/pinot/calcite/rel/rules/PinotRuleUtils.java
@@ -19,11 +19,13 @@
package org.apache.pinot.calcite.rel.rules;
import com.google.common.base.Preconditions;
+import com.google.common.collect.ImmutableList;
import java.util.ArrayList;
import java.util.EnumSet;
import java.util.List;
import javax.annotation.Nullable;
import org.apache.calcite.plan.Contexts;
+import org.apache.calcite.plan.RelTraitSet;
import org.apache.calcite.plan.hep.HepRelVertex;
import org.apache.calcite.rel.RelNode;
import org.apache.calcite.rel.core.Aggregate;
@@ -34,6 +36,8 @@ import org.apache.calcite.rel.core.Project;
import org.apache.calcite.rel.core.RelFactories;
import org.apache.calcite.rel.core.TableScan;
import org.apache.calcite.rel.core.Window;
+import org.apache.calcite.rel.logical.LogicalAsofJoin;
+import org.apache.calcite.rel.logical.LogicalJoin;
import org.apache.calcite.rex.RexCall;
import org.apache.calcite.rex.RexInputRef;
import org.apache.calcite.rex.RexLiteral;
@@ -302,4 +306,25 @@ public class PinotRuleUtils {
return null;
}
}
+
+ /// Returns `rel` with the given traits. `LogicalJoin#copy` and
`LogicalAsofJoin#copy` drop the trait set they are
+ /// given, so a copy of a join loses for example its distribution. This
creates them again with the traits instead.
+ public static RelNode withTraits(RelNode rel, RelTraitSet traitSet) {
+ if (rel.getTraitSet().equals(traitSet)) {
+ return rel;
+ }
+ if (rel instanceof LogicalJoin) {
+ LogicalJoin join = (LogicalJoin) rel;
+ return new LogicalJoin(join.getCluster(), traitSet, join.getHints(),
join.getLeft(), join.getRight(),
+ join.getCondition(), join.getVariablesSet(), join.getJoinType(),
join.isSemiJoinDone(),
+ ImmutableList.copyOf(join.getSystemFieldList()));
+ }
+ if (rel instanceof LogicalAsofJoin) {
+ LogicalAsofJoin join = (LogicalAsofJoin) rel;
+ return new LogicalAsofJoin(join.getCluster(), traitSet, join.getHints(),
join.getLeft(), join.getRight(),
+ join.getCondition(), join.getMatchCondition(), join.getJoinType(),
+ ImmutableList.copyOf(join.getSystemFieldList()));
+ }
+ return rel.copy(traitSet, rel.getInputs());
+ }
}
diff --git
a/pinot-query-planner/src/main/java/org/apache/pinot/calcite/rex/PinotSealedSearchOperator.java
b/pinot-query-planner/src/main/java/org/apache/pinot/calcite/rex/PinotSealedSearchOperator.java
new file mode 100644
index 00000000000..1a82c28ffc9
--- /dev/null
+++
b/pinot-query-planner/src/main/java/org/apache/pinot/calcite/rex/PinotSealedSearchOperator.java
@@ -0,0 +1,117 @@
+/**
+ * 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.calcite.rex;
+
+import java.util.function.Supplier;
+import javax.annotation.Nullable;
+import org.apache.calcite.plan.Strong;
+import org.apache.calcite.rel.type.RelDataTypeFactory;
+import org.apache.calcite.rex.RexLiteral;
+import org.apache.calcite.rex.RexUnknownAs;
+import org.apache.calcite.sql.SqlFunction;
+import org.apache.calcite.sql.SqlFunctionCategory;
+import org.apache.calcite.sql.SqlKind;
+import org.apache.calcite.sql.type.OperandTypes;
+import org.apache.calcite.sql.type.SqlReturnTypeInference;
+import org.apache.calcite.sql.type.SqlTypeName;
+import org.apache.calcite.util.Sarg;
+
+
+/// A `SEARCH(operand, sarg)` whose [Sarg] is hidden from Calcite while the
query is optimized.
+///
+/// Calcite keeps an IN list as one `SEARCH($x, Sarg[...])` call. Rules and
metadata handlers (for example predicate
+/// inference on joins) rebuild the Sarg's range set every time they simplify
a predicate list that holds it. For large
+/// lists this makes planning time grow with the list size times the number of
rule and metadata calls.
+///
+/// A call to this operator has exactly one operand, the operand of the
`SEARCH`. The Sarg literal is a field of the
+/// operator instance, not an operand. The call digest (for example
`$SEARCH#0($3)`), `equals` and `hashCode` therefore
+/// cost `O(1)`, and no Calcite code can read or rebuild the values.
[SearchSealer] creates one instance per distinct
+/// Sarg of a query and turns the calls back into `SEARCH` after optimization.
+///
+/// The operator keeps what Calcite knows about `SEARCH` without looking at
the values:
+/// - The kind is [SqlKind#OTHER_FUNCTION], not [SqlKind#SEARCH]: Calcite code
reads operand 1 of every `SEARCH` call
+/// as a Sarg literal.
+/// - The return type is `BOOLEAN`. It is nullable only if the operand is
nullable and the Sarg treats NULL as
+/// UNKNOWN. This is the rule of `SqlSearchOperator`.
+/// - The [Strong] policy is `ANY` for `NULL AS UNKNOWN` and `NOT_NULL`
otherwise. These are the answers that
+/// [Strong#isNull] gives for `SEARCH`, so a null-rejecting filter above an
outer join still makes it an inner join.
+/// - It is deterministic and safe (it never throws), like `SEARCH`. It is not
a dynamic function: Pinot does not
+/// relocate dynamic functions (see `PinotRuleUtils#isRelocatable`), and the
result only depends on the operand.
+/// A call with a literal operand is not folded during optimization. Pinot's
optimizer has no `RexExecutor`, so
+/// `ReduceExpressionsRule` does not evaluate calls, and
`PinotEvaluateLiteralRule` only evaluates registered scalar
+/// functions. `RexExpressionUtils` evaluates such a call after
optimization, as it does for `SEARCH`.
+/// - Equality is identity. [SearchSealer] shares one instance per distinct
Sarg within a query.
+///
+/// Instances are immutable and thread safe.
+public final class PinotSealedSearchOperator extends SqlFunction {
+ /// Name prefix of all sealed search operators.
+ static final String NAME_PREFIX = "$SEARCH#";
+
+ private final RexLiteral _sargLiteral;
+ private final Sarg<?> _sarg;
+
+ PinotSealedSearchOperator(int id, RexLiteral sargLiteral) {
+ this(id, sargLiteral, sargLiteral.getValueAs(Sarg.class));
+ }
+
+ private PinotSealedSearchOperator(int id, RexLiteral sargLiteral, Sarg<?>
sarg) {
+ super(NAME_PREFIX + id, SqlKind.OTHER_FUNCTION,
returnTypeInference(sarg.nullAs), null, OperandTypes.ANY,
+ SqlFunctionCategory.SYSTEM);
+ _sargLiteral = sargLiteral;
+ _sarg = sarg;
+ }
+
+ private static SqlReturnTypeInference returnTypeInference(RexUnknownAs
nullAs) {
+ return binding -> {
+ RelDataTypeFactory typeFactory = binding.getTypeFactory();
+ boolean nullable = nullAs == RexUnknownAs.UNKNOWN &&
binding.getOperandType(0).isNullable();
+ return
typeFactory.createTypeWithNullability(typeFactory.createSqlType(SqlTypeName.BOOLEAN),
nullable);
+ };
+ }
+
+ /// Returns the Sarg literal, which is operand 1 of the `SEARCH` call that
this operator seals.
+ public RexLiteral getSargLiteral() {
+ return _sargLiteral;
+ }
+
+ /// Returns the sealed Sarg.
+ Sarg<?> getSarg() {
+ return _sarg;
+ }
+
+ @Override
+ public Supplier<Strong.Policy> getStrongPolicyInference() {
+ return _sarg.nullAs == RexUnknownAs.UNKNOWN ? () -> Strong.Policy.ANY : ()
-> Strong.Policy.NOT_NULL;
+ }
+
+ @Override
+ public Boolean isSafeOperator() {
+ return true;
+ }
+
+ @Override
+ public boolean equals(@Nullable Object obj) {
+ return this == obj;
+ }
+
+ @Override
+ public int hashCode() {
+ return System.identityHashCode(this);
+ }
+}
diff --git
a/pinot-query-planner/src/main/java/org/apache/pinot/calcite/rex/SearchSealer.java
b/pinot-query-planner/src/main/java/org/apache/pinot/calcite/rex/SearchSealer.java
new file mode 100644
index 00000000000..1eaa006ca09
--- /dev/null
+++
b/pinot-query-planner/src/main/java/org/apache/pinot/calcite/rex/SearchSealer.java
@@ -0,0 +1,401 @@
+/**
+ * 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.calcite.rex;
+
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import javax.annotation.Nullable;
+import org.apache.calcite.plan.RelOptPredicateList;
+import org.apache.calcite.rel.RelHomogeneousShuttle;
+import org.apache.calcite.rel.RelNode;
+import org.apache.calcite.rel.core.Filter;
+import org.apache.calcite.rel.core.Join;
+import org.apache.calcite.rel.core.JoinRelType;
+import org.apache.calcite.rex.RexBuilder;
+import org.apache.calcite.rex.RexCall;
+import org.apache.calcite.rex.RexLiteral;
+import org.apache.calcite.rex.RexNode;
+import org.apache.calcite.rex.RexShuttle;
+import org.apache.calcite.rex.RexSimplify;
+import org.apache.calcite.rex.RexSubQuery;
+import org.apache.calcite.rex.RexVisitorImpl;
+import org.apache.calcite.sql.SqlBasicCall;
+import org.apache.calcite.sql.SqlCall;
+import org.apache.calcite.sql.SqlKind;
+import org.apache.calcite.sql.SqlNode;
+import org.apache.calcite.sql.SqlNodeList;
+import org.apache.calcite.sql.fun.SqlInOperator;
+import org.apache.calcite.sql.fun.SqlStdOperatorTable;
+import org.apache.calcite.sql.util.SqlBasicVisitor;
+import org.apache.calcite.util.Sarg;
+import org.apache.pinot.calcite.rel.rules.PinotRuleUtils;
+import org.apache.pinot.calcite.sql.fun.PinotInListOperator;
+
+
+/// Hides the large IN lists of one query from Calcite while the query is
optimized.
+///
+/// Calcite keeps an IN list as `SEARCH($x, Sarg[...])`. Many rules and
metadata handlers rebuild the Sarg's range set
+/// each time they touch it (`O(N log N)`), and each Calcite release and each
new planner rule adds such touches. This
+/// class replaces every `SEARCH` whose Sarg has at least [#getThreshold()]
ranges with a call to a
+/// [PinotSealedSearchOperator], which Calcite cannot look into, and turns it
back into the same `SEARCH` once
+/// optimization is done. Everything after optimization (EXPLAIN, plan node
conversion, broker pruning, the servers)
+/// sees the same `SEARCH` calls as without sealing.
+///
+/// Planning a query uses one instance, in this order:
+/// 1. [#markInLists] after validation. Large `IN` value lists skip
`SqlToRelConverter`'s expansion into `OR`, and
+/// `PinotConvertletTable` converts each marked call into one `SEARCH`.
Closing the [MarkedInLists] restores the SQL
+/// operators after the conversion.
+/// 2. [#seal(RelNode)] after `SqlToRelConverter`. It first simplifies each
filter and join condition that holds a
+/// large Sarg, like `RelBuilder#filter` does, so that the other predicates
on the same operand fold into the Sarg
+/// as they do without sealing (for example `x IS NOT NULL`, a range, or a
second list). Then it seals the large
+/// `SEARCH` calls, including the ones that Calcite built from a
user-written `x = 1 OR x = 2 OR ...`.
+/// 3. [#unseal(RelNode)] after the last optimizer program.
+///
+/// Sealed Sargs lose Calcite's value-level reasoning across plan nodes during
optimization. For example, a filter
+/// that a rule pushes onto another filter with a sealed list on the same
column is not merged with it. Filter
+/// push-down, transitive predicates across joins, outer join simplification
and all rules that treat the predicate
+/// as a whole work as before.
+///
+/// Null checks keep their meaning. Without null handling, the planner uses
Pinot's own `IS NULL` and `IS NOT NULL`
+/// operators (see `PinotOperatorTable`), which Calcite never merges into a
Sarg or drops. With null handling, Calcite
+/// can drop an `x IS NOT NULL` next to a sealed list on `x` as redundant, as
it does next to `x < 5`. This is correct,
+/// because the servers then apply SQL null semantics.
+///
+/// Not thread safe: a query is planned by a single thread.
+public final class SearchSealer {
+ private final int _threshold;
+ /// One operator per distinct Sarg literal (value and type), so that equal
lists stay equal (for example in both
+ /// branches of a common table expression).
+ private final Map<RexLiteral, PinotSealedSearchOperator> _operators = new
HashMap<>();
+
+ /// @param threshold the smallest number of Sarg ranges (IN list values) to
seal; 0 or less disables sealing
+ public SearchSealer(int threshold) {
+ _threshold = threshold;
+ }
+
+ /// Returns the smallest number of Sarg ranges (IN list values) that this
instance seals.
+ public int getThreshold() {
+ return _threshold;
+ }
+
+ /// Returns whether this instance seals anything.
+ public boolean isEnabled() {
+ return _threshold > 0;
+ }
+
+ /// Returns whether a Sarg is large enough to seal.
+ boolean shouldSeal(Sarg<?> sarg) {
+ // A Sarg that is all or none has at most one range and is never large.
Calcite has special cases for these
+ // (for example RexCall#isAlwaysTrue), so they always stay visible.
+ return _threshold > 0 && !sarg.isAll() && !sarg.isNone() &&
sarg.rangeSet.asRanges().size() >= _threshold;
+ }
+
+ /// Returns a sealed call for a `SEARCH` call whose Sarg is large enough, or
the call itself otherwise.
+ RexNode seal(RexBuilder rexBuilder, RexCall search) {
+ RexLiteral sargLiteral = (RexLiteral) search.getOperands().get(1);
+ Sarg<?> sarg = sargLiteral.getValueAs(Sarg.class);
+ if (sarg == null || !shouldSeal(sarg)) {
+ return search;
+ }
+ PinotSealedSearchOperator operator =
+ _operators.computeIfAbsent(sargLiteral, literal -> new
PinotSealedSearchOperator(_operators.size(), literal));
+ return rexBuilder.makeCall(search.getType(), operator,
List.of(search.getOperands().get(0)));
+ }
+
+ // --------------------------------------------------------------------------
+ // SQL level
+ // --------------------------------------------------------------------------
+
+ /// Sets [PinotInListOperator] on every `IN` and `NOT IN` call in the
validated tree whose value list has at least
+ /// [#getThreshold()] values. Close the result after `SqlToRelConverter` is
done, to restore the original operators:
+ ///
+ /// ```
+ /// try (SearchSealer.MarkedInLists ignored =
searchSealer.markInLists(validated)) {
+ /// relRoot = converter.convertQuery(validated, false, true);
+ /// }
+ /// ```
+ ///
+ /// Lists on a row (`(a, b) IN (...)`) and lists that contain a sub-query
keep the standard conversion.
+ public MarkedInLists markInLists(SqlNode validated) {
+ List<SqlBasicCall> marked = new ArrayList<>();
+ if (!isEnabled()) {
+ return new MarkedInLists(marked);
+ }
+ try {
+ validated.accept(new SqlBasicVisitor<Void>() {
+ @Override
+ public Void visit(SqlCall call) {
+ if (call instanceof SqlBasicCall && isLargeInList(call)) {
+ ((SqlBasicCall)
call).setOperator(PinotInListOperator.of(call.getKind()));
+ marked.add((SqlBasicCall) call);
+ // The values have no sub-queries (checked above), so only the
left operand needs a visit.
+ call.operand(0).accept(this);
+ return null;
+ }
+ return super.visit(call);
+ }
+ });
+ } catch (Throwable t) {
+ new MarkedInLists(marked).close();
+ throw t;
+ }
+ return new MarkedInLists(marked);
+ }
+
+ /// The `IN` and `NOT IN` calls that [#markInLists] marked. [#close()]
restores their original operators.
+ public static final class MarkedInLists implements AutoCloseable {
+ private final List<SqlBasicCall> _calls;
+
+ private MarkedInLists(List<SqlBasicCall> calls) {
+ _calls = calls;
+ }
+
+ @Override
+ public void close() {
+ for (SqlBasicCall call : _calls) {
+ call.setOperator(((PinotInListOperator)
call.getOperator()).getOriginal());
+ }
+ }
+ }
+
+ private boolean isLargeInList(SqlCall call) {
+ SqlKind kind = call.getKind();
+ if ((kind != SqlKind.IN && kind != SqlKind.NOT_IN) || !(call.getOperator()
instanceof SqlInOperator)
+ || call.operandCount() != 2) {
+ return false;
+ }
+ if (!(call.operand(1) instanceof SqlNodeList) || call.operand(0).getKind()
== SqlKind.ROW) {
+ return false;
+ }
+ SqlNodeList values = call.operand(1);
+ return values.size() >= _threshold && !containsQueryOrRow(values);
+ }
+
+ private static boolean containsQueryOrRow(SqlNodeList values) {
+ SqlBasicVisitor<Boolean> finder = new SqlBasicVisitor<>() {
+ @Override
+ public Boolean visit(SqlCall call) {
+ SqlKind kind = call.getKind();
+ if (kind == SqlKind.ROW || kind == SqlKind.SCALAR_QUERY ||
kind.belongsTo(SqlKind.QUERY)) {
+ return true;
+ }
+ for (SqlNode operand : call.getOperandList()) {
+ if (operand != null && Boolean.TRUE.equals(operand.accept(this))) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ @Override
+ public Boolean visit(SqlNodeList nodeList) {
+ for (SqlNode node : nodeList) {
+ if (node != null && Boolean.TRUE.equals(node.accept(this))) {
+ return true;
+ }
+ }
+ return false;
+ }
+ };
+ return Boolean.TRUE.equals(values.accept(finder));
+ }
+
+ // --------------------------------------------------------------------------
+ // Relational level
+ // --------------------------------------------------------------------------
+
+ /// Seals every large `SEARCH` in a plan that `SqlToRelConverter` produced.
+ ///
+ /// A filter or join condition that holds a large Sarg, or a large `AND` /
`OR` of comparisons, is simplified first,
+ /// the same way `RelBuilder#filter` does it. Field trimming usually did
this already when it rebuilt the filter,
+ /// but not when the filter keeps all its fields. This folds `x = 1 OR x = 2
OR ...` into a Sarg, and it folds the
+ /// other predicates on the same operand into the Sarg (for example `x IS
NOT NULL` into `NULL AS FALSE`), exactly
+ /// as the optimizer would do without sealing. Once sealed, such a predicate
cannot merge anymore, and Calcite could
+ /// drop an `x IS NOT NULL` next to it as redundant.
+ public RelNode seal(RelNode rel) {
+ return isEnabled() ? rel.accept(new RelSealer()) : rel;
+ }
+
+ /// Restores the `SEARCH` calls that this instance sealed.
+ public RelNode unseal(RelNode rel) {
+ return _operators.isEmpty() ? rel : rel.accept(new RelUnsealer());
+ }
+
+ /// Returns the `SEARCH` call that a sealed call stands for.
+ static RexNode unsealCall(RexBuilder rexBuilder, RexCall sealed) {
+ PinotSealedSearchOperator operator = (PinotSealedSearchOperator)
sealed.getOperator();
+ return rexBuilder.makeCall(sealed.getType(), SqlStdOperatorTable.SEARCH,
+ List.of(sealed.getOperands().get(0), operator.getSargLiteral()));
+ }
+
+ private final class RelSealer extends RelHomogeneousShuttle {
+ @Override
+ public RelNode visit(RelNode other) {
+ RelNode rel = foldLargeLists(super.visit(other));
+ return rel.accept(new RexSealer(rel.getCluster().getRexBuilder()));
+ }
+
+ private RelNode foldLargeLists(RelNode rel) {
+ if (rel instanceof Filter) {
+ Filter filter = (Filter) rel;
+ if (hasLargeList(filter.getCondition())) {
+ return filter.copy(filter.getTraitSet(), filter.getInput(),
+ simplifyCondition(filter.getCluster().getRexBuilder(),
filter.getCondition()));
+ }
+ } else if (rel instanceof Join) {
+ Join join = (Join) rel;
+ if (join.getJoinType() != JoinRelType.LEFT_MARK &&
hasLargeList(join.getCondition())) {
+ return join.copy(join.getTraitSet(),
simplifyCondition(join.getCluster().getRexBuilder(),
+ join.getCondition()), join.getLeft(), join.getRight(),
join.getJoinType(), join.isSemiJoinDone());
+ }
+ }
+ return rel;
+ }
+ }
+
+ /// Returns whether a condition has a `SEARCH` that is large enough to seal,
or an AND or OR with so many comparisons
+ /// of one operand to literals that Calcite could fold them into such a
Sarg. An AND of `x <> v` comparisons folds
+ /// into `n + 1` ranges. Comparisons of different operands cannot fold into
one Sarg, so they are counted apart.
+ private boolean hasLargeList(RexNode condition) {
+ int minComparisons = Math.max(_threshold - 1, 2);
+ return Boolean.TRUE.equals(condition.accept(new
RexVisitorImpl<Boolean>(true) {
+ @Override
+ public Boolean visitCall(RexCall call) {
+ SqlKind kind = call.getKind();
+ if (kind == SqlKind.SEARCH) {
+ Sarg<?> sarg = ((RexLiteral)
call.getOperands().get(1)).getValueAs(Sarg.class);
+ if (sarg != null && shouldSeal(sarg)) {
+ return true;
+ }
+ }
+ if ((kind == SqlKind.AND || kind == SqlKind.OR) &&
call.getOperands().size() >= minComparisons) {
+ Map<RexNode, Integer> comparisons = new HashMap<>();
+ for (RexNode operand : call.getOperands()) {
+ RexNode compared = comparedToLiteral(operand);
+ if (compared != null && comparisons.merge(compared, 1,
Integer::sum) >= minComparisons) {
+ return true;
+ }
+ }
+ }
+ for (RexNode operand : call.getOperands()) {
+ if (Boolean.TRUE.equals(operand.accept(this))) {
+ return true;
+ }
+ }
+ return false;
+ }
+ }));
+ }
+
+ /// Returns the operand that a comparison to a literal compares, or null if
the node is not such a comparison.
+ @Nullable
+ private static RexNode comparedToLiteral(RexNode node) {
+ if (!node.isA(SqlKind.COMPARISON) || !(node instanceof RexCall)) {
+ return null;
+ }
+ List<RexNode> operands = ((RexCall) node).getOperands();
+ if (operands.size() != 2) {
+ return null;
+ }
+ if (operands.get(1) instanceof RexLiteral) {
+ return operands.get(0);
+ }
+ return operands.get(0) instanceof RexLiteral ? operands.get(1) : null;
+ }
+
+ /// Simplifies a filter or join condition like `RelBuilder#filter` does.
Unlike `RexSimplify#simplifyUnknownAsFalse`,
+ /// this does not simplify each OR term against the other terms, which is
quadratic in the number of terms.
+ private static RexNode simplifyCondition(RexBuilder rexBuilder, RexNode
condition) {
+ RexNode simplified = new RexSimplify(rexBuilder,
RelOptPredicateList.EMPTY, PinotRexExecutor.INSTANCE)
+ .simplifyFilterPredicates(List.of(condition));
+ return simplified != null ? simplified : rexBuilder.makeLiteral(false);
+ }
+
+ private final class RexSealer extends RexShuttle {
+ private final RexBuilder _rexBuilder;
+
+ RexSealer(RexBuilder rexBuilder) {
+ _rexBuilder = rexBuilder;
+ }
+
+ @Override
+ public RexNode visitCall(RexCall call) {
+ RexNode visited = super.visitCall(call);
+ if (visited.getKind() == SqlKind.SEARCH) {
+ return seal(_rexBuilder, (RexCall) visited);
+ }
+ return visited;
+ }
+
+ @Override
+ public RexNode visitSubQuery(RexSubQuery subQuery) {
+ RexSubQuery visited = (RexSubQuery) super.visitSubQuery(subQuery);
+ RelNode rel = visited.rel.accept(new RelSealer());
+ return rel == visited.rel ? visited : visited.clone(rel);
+ }
+ }
+
+ private static final class RelUnsealer extends RelHomogeneousShuttle {
+ @Override
+ public RelNode visit(RelNode other) {
+ RelNode rel = super.visit(other);
+ rel = rel.accept(new RexUnsealer(rel.getCluster().getRexBuilder()));
+ // Un-sealing runs after the trait program: a copy must keep the traits
of the original, for example its
+ // distribution.
+ return rel == other ? other : PinotRuleUtils.withTraits(rel,
other.getTraitSet());
+ }
+ }
+
+
+ private static final class RexUnsealer extends RexShuttle {
+ private final RexBuilder _rexBuilder;
+
+ RexUnsealer(RexBuilder rexBuilder) {
+ _rexBuilder = rexBuilder;
+ }
+
+ @Override
+ public RexNode visitCall(RexCall call) {
+ if (call.getKind() == SqlKind.NOT &&
isSealed(call.getOperands().get(0))) {
+ // RexSimplify pushes NOT into a SEARCH by negating its Sarg, which it
could not do while the Sarg was sealed.
+ RexCall sealed = (RexCall) call.getOperands().get(0);
+ PinotSealedSearchOperator operator = (PinotSealedSearchOperator)
sealed.getOperator();
+ RexLiteral negated =
+ _rexBuilder.makeSearchArgumentLiteral(operator.getSarg().negate(),
operator.getSargLiteral().getType());
+ return _rexBuilder.makeCall(call.getType(), SqlStdOperatorTable.SEARCH,
+ List.of(sealed.getOperands().get(0).accept(this), negated));
+ }
+ RexNode visited = super.visitCall(call);
+ return isSealed(visited) ? unsealCall(_rexBuilder, (RexCall) visited) :
visited;
+ }
+
+ private boolean isSealed(RexNode node) {
+ return node instanceof RexCall && ((RexCall) node).getOperator()
instanceof PinotSealedSearchOperator;
+ }
+
+ @Override
+ public RexNode visitSubQuery(RexSubQuery subQuery) {
+ RexSubQuery visited = (RexSubQuery) super.visitSubQuery(subQuery);
+ RelNode rel = visited.rel.accept(new RelUnsealer());
+ return rel == visited.rel ? visited : visited.clone(rel);
+ }
+ }
+}
diff --git
a/pinot-query-planner/src/main/java/org/apache/pinot/calcite/sql/fun/PinotInListOperator.java
b/pinot-query-planner/src/main/java/org/apache/pinot/calcite/sql/fun/PinotInListOperator.java
new file mode 100644
index 00000000000..3143f763838
--- /dev/null
+++
b/pinot-query-planner/src/main/java/org/apache/pinot/calcite/sql/fun/PinotInListOperator.java
@@ -0,0 +1,93 @@
+/**
+ * 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.calcite.sql.fun;
+
+import org.apache.calcite.rel.type.RelDataType;
+import org.apache.calcite.sql.SqlBinaryOperator;
+import org.apache.calcite.sql.SqlCall;
+import org.apache.calcite.sql.SqlKind;
+import org.apache.calcite.sql.SqlSpecialOperator;
+import org.apache.calcite.sql.SqlWriter;
+import org.apache.calcite.sql.fun.SqlStdOperatorTable;
+import org.apache.calcite.sql.type.InferTypes;
+import org.apache.calcite.sql.type.ReturnTypes;
+import org.apache.calcite.sql.validate.SqlValidator;
+import org.apache.calcite.sql.validate.SqlValidatorScope;
+
+
+/// Marks an `IN` or `NOT IN` call with a large value list while the query is
converted to relational algebra.
+///
+/// `SqlToRelConverter` registers every `IN` call as a sub-query and, because
Pinot sets `inSubQueryThreshold` to
+/// `Integer.MAX_VALUE`, expands a value list into `x = v1 OR x = v2 OR ...`.
Simplifying this OR costs `O(N log N)` in
+/// a WHERE clause and much more elsewhere (CASE, SELECT list, aggregate
FILTER), because `RexSimplify` then simplifies
+/// every term against the negation of all earlier terms.
+///
+/// After validation, `SearchSealer#markInLists` sets this operator on large
`IN` calls in place. This keeps the node
+/// identity, so the types that the validator derived stay valid. Its kind is
not `IN`, so `SqlToRelConverter` does not
+/// expand the list. It converts the call with `PinotConvertletTable`, which
builds one `SEARCH`. The original
+/// operator is set again after the conversion, so the rest of Pinot never
sees this operator.
+///
+/// The name, precedence, return type and unparse are the ones of the original
operator. `SqlNode#equalsDeep` compares
+/// operators by name, so a marked expression still matches the same
expression in the GROUP BY clause.
+///
+/// The two instances are immutable and thread safe.
+public final class PinotInListOperator extends SqlSpecialOperator {
+ public static final PinotInListOperator IN = new
PinotInListOperator(SqlStdOperatorTable.IN);
+ public static final PinotInListOperator NOT_IN = new
PinotInListOperator(SqlStdOperatorTable.NOT_IN);
+
+ private final SqlBinaryOperator _original;
+
+ private PinotInListOperator(SqlBinaryOperator original) {
+ super(original.getName(), SqlKind.OTHER_FUNCTION, original.getLeftPrec(),
true, ReturnTypes.BOOLEAN_NULLABLE,
+ InferTypes.FIRST_KNOWN, null);
+ _original = original;
+ }
+
+ /// Returns the marker for an `IN` or `NOT IN` call.
+ public static PinotInListOperator of(SqlKind kind) {
+ switch (kind) {
+ case IN:
+ return IN;
+ case NOT_IN:
+ return NOT_IN;
+ default:
+ throw new IllegalArgumentException("Not an IN list kind: " + kind);
+ }
+ }
+
+ /// Returns whether this operator marks a `NOT IN` call.
+ public boolean isNegated() {
+ return this == NOT_IN;
+ }
+
+ /// Returns the operator that this operator replaces.
+ public SqlBinaryOperator getOriginal() {
+ return _original;
+ }
+
+ @Override
+ public RelDataType deriveType(SqlValidator validator, SqlValidatorScope
scope, SqlCall call) {
+ return _original.deriveType(validator, scope, call);
+ }
+
+ @Override
+ public void unparse(SqlWriter writer, SqlCall call, int leftPrec, int
rightPrec) {
+ _original.unparse(writer, call, leftPrec, rightPrec);
+ }
+}
diff --git
a/pinot-query-planner/src/main/java/org/apache/pinot/calcite/sql2rel/PinotConvertletTable.java
b/pinot-query-planner/src/main/java/org/apache/pinot/calcite/sql2rel/PinotConvertletTable.java
index e908fc30b9e..124760d4987 100644
---
a/pinot-query-planner/src/main/java/org/apache/pinot/calcite/sql2rel/PinotConvertletTable.java
+++
b/pinot-query-planner/src/main/java/org/apache/pinot/calcite/sql2rel/PinotConvertletTable.java
@@ -18,13 +18,24 @@
*/
package org.apache.pinot.calcite.sql2rel;
+import java.util.ArrayList;
+import java.util.HashMap;
import java.util.List;
+import java.util.Map;
+import java.util.Objects;
import javax.annotation.Nullable;
+import org.apache.calcite.plan.RelOptPredicateList;
import org.apache.calcite.rex.RexBuilder;
+import org.apache.calcite.rex.RexCall;
+import org.apache.calcite.rex.RexLiteral;
import org.apache.calcite.rex.RexNode;
+import org.apache.calcite.rex.RexSimplify;
+import org.apache.calcite.rex.RexUnknownAs;
+import org.apache.calcite.rex.RexUtil;
import org.apache.calcite.sql.SqlCall;
import org.apache.calcite.sql.SqlKind;
import org.apache.calcite.sql.SqlNode;
+import org.apache.calcite.sql.SqlNodeList;
import org.apache.calcite.sql.fun.SqlBetweenOperator;
import org.apache.calcite.sql.fun.SqlStdOperatorTable;
import org.apache.calcite.sql2rel.SqlRexContext;
@@ -32,10 +43,14 @@ import org.apache.calcite.sql2rel.SqlRexConvertlet;
import org.apache.calcite.sql2rel.SqlRexConvertletTable;
import org.apache.calcite.sql2rel.StandardConvertletTable;
import org.apache.calcite.util.Litmus;
+import org.apache.calcite.util.Sarg;
+import org.apache.pinot.calcite.rex.PinotRexExecutor;
+import org.apache.pinot.calcite.sql.fun.PinotInListOperator;
/// PinotConvertletTable is a wrapper of [StandardConvertletTable] with the
customizations of not converting
-/// certain SqlCalls, e.g. TIMESTAMPADD, TIMESTAMPDIFF.
+/// certain SqlCalls, e.g. TIMESTAMPADD, TIMESTAMPDIFF. It also converts the
large IN lists that `SearchSealer` marked
+/// into one `SEARCH` each (see [InListConvertlet]).
public class PinotConvertletTable implements SqlRexConvertletTable {
public static final PinotConvertletTable INSTANCE = new
PinotConvertletTable();
@@ -53,6 +68,9 @@ public class PinotConvertletTable implements
SqlRexConvertletTable {
@Nullable
@Override
public SqlRexConvertlet get(SqlCall call) {
+ if (call.getOperator() instanceof PinotInListOperator) {
+ return InListConvertlet.INSTANCE;
+ }
switch (call.getKind()) {
case TIMESTAMP_ADD:
return TimestampAddConvertlet.INSTANCE;
@@ -143,6 +161,126 @@ public class PinotConvertletTable implements
SqlRexConvertletTable {
}
}
+ /// Converts an `IN` or `NOT IN` call that `SearchSealer#markInLists` marked
with [PinotInListOperator].
+ ///
+ /// The result is the expression that `SqlToRelConverter` and a later
`RexSimplify` would build, without building
+ /// and simplifying an `OR` of `N` terms (which is cubic in `N` outside a
filter):
+ /// 1. Like `SqlToRelConverter#convertInToOr`, each value becomes `operand =
value`, converted by the standard
+ /// convertlets. So operand and literal types (including casts from type
coercion) are the same as without
+ /// sealing. Each comparison is then simplified alone, which reduces
casts of literals such as `CAST('1' AS
+ /// INTEGER)` (with the query's cached cast executables) and turns `x =
NULL` into `NULL`.
+ /// 2. The comparisons of one operand with non-null literals become one
`SEARCH` ([RexBuilder#makeIn]). For `NOT IN`
+ /// the Sarg is negated, as `RexSimplify` does for `NOT(SEARCH)`.
`SearchSealer#seal(RelNode)` seals the
+ /// `SEARCH` after `SqlToRelConverter`, once Calcite has folded the other
predicates on the same operand (for
+ /// example `x IS NOT NULL`) into it, as it does without sealing.
+ /// 3. Other comparisons (NULL values, values that are not literals) stay as
they are: `IN` is the `OR` of all terms
+ /// and `NOT IN` is the `AND` of their negations. This is exact
three-valued logic: for example,
+ /// `x IN (1, 2, NULL)` becomes `OR(SEARCH(x, [1, 2]), NULL)`.
+ private static class InListConvertlet implements SqlRexConvertlet {
+ static final InListConvertlet INSTANCE = new InListConvertlet();
+
+ @Override
+ public RexNode convertCall(SqlRexContext cx, SqlCall call) {
+ RexBuilder rexBuilder = cx.getRexBuilder();
+ boolean negated = ((PinotInListOperator) call.getOperator()).isNegated();
+ SqlNode operand = call.operand(0);
+ SqlNodeList values = call.operand(1);
+ RexSimplify simplify = new RexSimplify(rexBuilder,
RelOptPredicateList.EMPTY, PinotRexExecutor.INSTANCE);
+
+ // Each slot is a term, or the points of one operand at the position of
its first point (like RexSimplify's Sarg
+ // collector), so the terms keep the order that Calcite gives them.
+ List<Object> slots = new ArrayList<>();
+ Map<RexNode, Points> pointsByOperand = new HashMap<>();
+ for (SqlNode value : values) {
+ SqlCall equals =
SqlStdOperatorTable.EQUALS.createCall(value.getParserPosition(), operand,
value);
+ RexNode term =
simplify.simplifyUnknownAs(cx.convertExpression(equals), RexUnknownAs.UNKNOWN);
+ RexNode pointOperand = getPointOperand(term);
+ if (pointOperand == null) {
+ slots.add(term);
+ continue;
+ }
+ Points points = pointsByOperand.get(pointOperand);
+ if (points == null) {
+ points = new Points(pointOperand, new ArrayList<>());
+ pointsByOperand.put(pointOperand, points);
+ slots.add(points);
+ }
+ points.literals().add(getOther(term, pointOperand));
+ }
+
+ List<RexNode> terms = new ArrayList<>(slots.size());
+ for (Object slot : slots) {
+ if (slot instanceof Points points) {
+ RexNode in = rexBuilder.makeIn(points.operand(), points.literals());
+ terms.add(negated ? negate(rexBuilder, in) : in);
+ } else {
+ terms.add(negated ? negate(rexBuilder, (RexNode) slot) : (RexNode)
slot);
+ }
+ }
+ RexNode result;
+ if (negated) {
+ result = RexUtil.composeConjunction(rexBuilder, terms);
+ } else {
+ result = RexUtil.composeDisjunction(rexBuilder, terms);
+ }
+ return StandardConvertletTable.castToValidatedType(call, result,
cx.getValidator(), rexBuilder);
+ }
+
+ /// Returns the operand of `term` if it compares a deterministic
expression with a non-null literal (the terms
+ /// that `RexSimplify`'s Sarg collector accepts for `=`), or null
otherwise.
+ @Nullable
+ private static RexNode getPointOperand(RexNode term) {
+ if (term.getKind() != SqlKind.EQUALS) {
+ return null;
+ }
+ List<RexNode> operands = ((RexCall) term).getOperands();
+ RexNode left = operands.get(0);
+ RexNode right = operands.get(1);
+ if (isNonNullLiteral(right) && RexUtil.isDeterministic(left)) {
+ return left;
+ }
+ if (isNonNullLiteral(left) && RexUtil.isDeterministic(right)) {
+ return right;
+ }
+ return null;
+ }
+
+ private static RexNode getOther(RexNode term, RexNode pointOperand) {
+ List<RexNode> operands = ((RexCall) term).getOperands();
+ return operands.get(0) == pointOperand ? operands.get(1) :
operands.get(0);
+ }
+
+ private static boolean isNonNullLiteral(RexNode node) {
+ return node instanceof RexLiteral && !((RexLiteral) node).isNull();
+ }
+
+ private static RexNode negate(RexBuilder rexBuilder, RexNode node) {
+ if (node.getKind() == SqlKind.SEARCH) {
+ RexCall search = (RexCall) node;
+ RexLiteral sargLiteral = (RexLiteral) search.getOperands().get(1);
+ Sarg<?> sarg =
Objects.requireNonNull(sargLiteral.getValueAs(Sarg.class));
+ return rexBuilder.makeCall(SqlStdOperatorTable.SEARCH,
search.getOperands().get(0),
+ rexBuilder.makeSearchArgumentLiteral(sarg.negate(),
sargLiteral.getType()));
+ }
+ if (RexUtil.isNullLiteral(node, false)) {
+ // NOT(NULL) is NULL.
+ return node;
+ }
+ if (node instanceof RexCall) {
+ // For example, x = y becomes x <> y.
+ RexNode negated = RexUtil.negate(rexBuilder, (RexCall) node);
+ if (negated != null) {
+ return negated;
+ }
+ }
+ return RexUtil.not(node);
+ }
+ }
+
+ /// The literals that one operand of an IN list is compared with.
+ private record Points(RexNode operand, List<RexNode> literals) {
+ }
+
/// Check if a comparison call involves ROW expressions.
private static boolean isRowComparison(SqlCall call) {
if (call.getOperandList().size() != 2) {
diff --git
a/pinot-query-planner/src/main/java/org/apache/pinot/query/QueryEnvironment.java
b/pinot-query-planner/src/main/java/org/apache/pinot/query/QueryEnvironment.java
index d07f2089082..34cc9e7c92d 100644
---
a/pinot-query-planner/src/main/java/org/apache/pinot/query/QueryEnvironment.java
+++
b/pinot-query-planner/src/main/java/org/apache/pinot/query/QueryEnvironment.java
@@ -64,6 +64,7 @@ import
org.apache.pinot.calcite.rel.rules.PinotRelDistributionTraitRule;
import org.apache.pinot.calcite.rel.rules.PinotRuleUtils;
import org.apache.pinot.calcite.rel.rules.PinotSortExchangeCopyRule;
import org.apache.pinot.calcite.rex.PinotRexExecutor;
+import org.apache.pinot.calcite.rex.SearchSealer;
import org.apache.pinot.calcite.sql.fun.PinotOperatorTable;
import org.apache.pinot.calcite.sql2rel.PinotConvertletTable;
import org.apache.pinot.calcite.sql2rel.PinotRelDecorrelator;
@@ -458,11 +459,13 @@ public class QueryEnvironment {
try {
RexBuilder rexBuilder = new RexBuilder(_typeFactory);
RelOptCluster cluster =
RelOptCluster.create(plannerContext.getRelOptPlanner(), rexBuilder);
+ SearchSealer searchSealer = plannerContext.getSearchSealer();
SqlToRelConverter converter =
new SqlToRelConverter(plannerContext.getPlanner(),
plannerContext.getValidator(), _catalogReader, cluster,
PinotConvertletTable.INSTANCE,
_config.getSqlToRelConverterConfig());
RelRoot relRoot;
- try {
+ // Large IN lists skip SqlToRelConverter's expansion into OR;
PinotConvertletTable builds one SEARCH for each.
+ try (SearchSealer.MarkedInLists ignored =
searchSealer.markInLists(sqlNode)) {
relRoot = converter.convertQuery(sqlNode, false, true);
} catch (Throwable e) {
throw new RuntimeException("Failed to convert query to relational
expression:\n" + sqlNode, e);
@@ -484,6 +487,8 @@ public class QueryEnvironment {
} catch (Throwable e) {
throw new RuntimeException("Failed to trim unused fields from
query:\n" + RelOptUtil.toString(rootNode), e);
}
+ // Hide the large SEARCH calls from the optimizer. SearchSealer#unseal
restores them at the end of optimize().
+ rootNode = searchSealer.seal(rootNode);
return relRoot.withRel(rootNode);
} catch (QueryException e) {
throw e;
@@ -517,7 +522,8 @@ public class QueryEnvironment {
listener.populateRuleTimings();
RelOptPlanner traitPlanner = plannerContext.getRelTraitPlanner();
traitPlanner.setRoot(optimized);
- return traitPlanner.findBestExp();
+ // Everything after optimization (EXPLAIN, plan node conversion,
physical planning) sees plain SEARCH calls.
+ return
plannerContext.getSearchSealer().unseal(traitPlanner.findBestExp());
} catch (Throwable e) {
throw QueryErrorCode.QUERY_PLANNING.asException("Error optimizing query:
" + e.getMessage(), e);
}
@@ -939,6 +945,13 @@ public class QueryEnvironment {
default int defaultSortExchangeCopyLimit() {
return
PinotSortExchangeCopyRule.SORT_EXCHANGE_COPY.config.getFetchLimitThreshold();
}
+
+ /// See [CommonConstants.Broker#CONFIG_OF_SEALED_IN_LIST_THRESHOLD]. Can
be overridden per query with
+ ///
[CommonConstants.Broker.Request.QueryOptionKey#SEALED_IN_LIST_THRESHOLD].
+ @Value.Default
+ default int defaultSealedInListThreshold() {
+ return CommonConstants.Broker.DEFAULT_SEALED_IN_LIST_THRESHOLD;
+ }
}
/// A query that have been parsed, validates, transformed into a [RelNode]
and optimized with Calcite.
diff --git
a/pinot-query-planner/src/main/java/org/apache/pinot/query/context/PlannerContext.java
b/pinot-query-planner/src/main/java/org/apache/pinot/query/context/PlannerContext.java
index f5748199acb..c5708d3cddf 100644
---
a/pinot-query-planner/src/main/java/org/apache/pinot/query/context/PlannerContext.java
+++
b/pinot-query-planner/src/main/java/org/apache/pinot/query/context/PlannerContext.java
@@ -34,6 +34,8 @@ import org.apache.calcite.rel.type.RelDataTypeFactory;
import org.apache.calcite.sql.SqlExplainFormat;
import org.apache.calcite.sql.validate.SqlValidator;
import org.apache.calcite.tools.FrameworkConfig;
+import org.apache.pinot.calcite.rex.SearchSealer;
+import org.apache.pinot.common.utils.config.QueryOptionsUtils;
import org.apache.pinot.query.QueryEnvironment;
import org.apache.pinot.query.planner.logical.LogicalPlanner;
import org.apache.pinot.query.validate.Validator;
@@ -61,6 +63,7 @@ public class PlannerContext implements AutoCloseable, Context
{
private final SqlExplainFormat _sqlExplainFormat;
@Nullable
private final PhysicalPlannerContext _physicalPlannerContext;
+ private final SearchSealer _searchSealer;
/// Set by the approximate aggregation rewrite rule when it rewrites at
least one aggregation, so that the broker
/// can report on the response that the results are approximate. Written
during planning, which is single threaded
/// for a query, and read after planning completes.
@@ -79,6 +82,7 @@ public class PlannerContext implements AutoCloseable, Context
{
_plannerOutput = new HashMap<>();
_sqlExplainFormat = sqlExplainFormat;
_physicalPlannerContext = physicalPlannerContext;
+ _searchSealer = createSearchSealer(options, envConfig);
}
/// Test factory: creates a minimal [PlannerContext] without going through
@@ -101,6 +105,12 @@ public class PlannerContext implements AutoCloseable,
Context {
_plannerOutput = new HashMap<>();
_sqlExplainFormat = null;
_physicalPlannerContext = null;
+ _searchSealer = createSearchSealer(options, envConfig);
+ }
+
+ private static SearchSealer createSearchSealer(Map<String, String> options,
QueryEnvironment.Config envConfig) {
+ return new SearchSealer(
+ QueryOptionsUtils.getSealedInListThreshold(options,
envConfig.defaultSealedInListThreshold()));
}
public PlannerImpl getPlanner() {
@@ -172,6 +182,11 @@ public class PlannerContext implements AutoCloseable,
Context {
return _physicalPlannerContext;
}
+ /// Returns the query's [SearchSealer], which hides large IN lists from
Calcite during optimization.
+ public SearchSealer getSearchSealer() {
+ return _searchSealer;
+ }
+
public boolean isUsePhysicalOptimizer() {
return _physicalPlannerContext != null;
}
diff --git
a/pinot-query-planner/src/main/java/org/apache/pinot/query/planner/logical/RexExpressionUtils.java
b/pinot-query-planner/src/main/java/org/apache/pinot/query/planner/logical/RexExpressionUtils.java
index 4458d4ed23a..73c6ca8bb5f 100644
---
a/pinot-query-planner/src/main/java/org/apache/pinot/query/planner/logical/RexExpressionUtils.java
+++
b/pinot-query-planner/src/main/java/org/apache/pinot/query/planner/logical/RexExpressionUtils.java
@@ -20,10 +20,13 @@ package org.apache.pinot.query.planner.logical;
import com.google.common.base.Preconditions;
import com.google.common.collect.BoundType;
+import com.google.common.collect.ImmutableRangeSet;
import com.google.common.collect.Range;
+import com.google.common.collect.RangeSet;
import java.math.BigDecimal;
import java.util.ArrayList;
import java.util.Calendar;
+import java.util.Collection;
import java.util.List;
import java.util.Set;
import java.util.UUID;
@@ -50,6 +53,7 @@ import org.apache.calcite.tools.RelBuilder;
import org.apache.calcite.util.NlsString;
import org.apache.calcite.util.Sarg;
import org.apache.calcite.util.TimestampString;
+import org.apache.pinot.calcite.rex.PinotSealedSearchOperator;
import org.apache.pinot.common.function.scalar.arithmetic.NegateScalarFunction;
import org.apache.pinot.common.utils.DataSchema.ColumnDataType;
import org.apache.pinot.spi.utils.BooleanUtils;
@@ -279,13 +283,19 @@ public class RexExpressionUtils {
}
public static RexExpression fromRexCall(RexCall rexCall) {
+ if (rexCall.op instanceof PinotSealedSearchOperator) {
+ // Sealed searches are restored to SEARCH at the end of optimization.
Convert them the same way in case a
+ // conversion runs before that.
+ return handleSearch(rexCall.operands.get(0),
((PinotSealedSearchOperator) rexCall.op).getSargLiteral());
+ }
switch (rexCall.op.kind) {
case CAST:
return handleCast(rexCall);
case REINTERPRET:
return handleReinterpret(rexCall);
case SEARCH:
- return handleSearch(rexCall);
+ assert rexCall.operands.size() == 2;
+ return handleSearch(rexCall.operands.get(0), (RexLiteral)
rexCall.operands.get(1));
case MINUS_PREFIX:
// Without this explicit case the default branch calls
getFunctionName(), which returns
// SqlKind.MINUS_PREFIX.name() = "MINUS_PREFIX". That canonicalizes to
"minusprefix", which is
@@ -346,10 +356,7 @@ public class RexExpressionUtils {
return fromRexNode(rexCall.operands.get(0));
}
- private static RexExpression handleSearch(RexCall rexCall) {
- assert rexCall.operands.size() == 2;
- RexNode leftOperand = rexCall.operands.get(0);
- RexLiteral searchArgument = (RexLiteral) rexCall.operands.get(1);
+ private static RexExpression handleSearch(RexNode leftOperand, RexLiteral
searchArgument) {
ColumnDataType dataType =
RelToPlanNodeConverter.convertToColumnDataType(searchArgument.getType());
Sarg sarg = searchArgument.getValueAs(Sarg.class);
assert sarg != null;
@@ -371,11 +378,89 @@ public class RexExpressionUtils {
if (leftOperand instanceof RexLiteral) {
return evaluateLiteralOrRanges((RexLiteral) leftOperand,
sarg.rangeSet.asRanges(), sarg.nullAs);
}
- RexExpression orExpr = convertRangesToOr(dataType, leftOperand,
sarg.rangeSet.asRanges());
- return addNullCheckIfRequired(leftOperand, sarg.nullAs, orExpr);
+ RexExpression rangesExpr = convertRanges(dataType, leftOperand,
sarg.rangeSet);
+ return addNullCheckIfRequired(leftOperand, sarg.nullAs, rangesExpr);
+ }
+ }
+
+ /// The smallest number of single values in a range set with other ranges
that ships as one `IN` or `NOT_IN`.
+ /// Below it, one range per value costs the servers little.
+ static final int MIN_POINTS_FOR_IN = 20;
+
+ /// Converts a range set that is neither only points nor only the complement
of points.
+ ///
+ /// Each range becomes comparisons, except when the range set holds many
single points. A large OR of
+ /// `x >= v AND x <= v` terms costs the servers one scan per value, so:
+ /// - When the ranges hold at least [#MIN_POINTS_FOR_IN] points, for example
`x IN (<list>) OR x > 10`, the points
+ /// become one `IN` and only the other ranges become comparisons:
`OR(IN(x, <list>), x > 10)`.
+ /// - When the complement holds at least [#MIN_POINTS_FOR_IN] points, for
example `x NOT IN (<list>) AND x > 0`, the
+ /// result is `AND(NOT_IN(x, <list>), x > 0)`.
+ /// - When both forms apply, the one with fewer values is used.
+ ///
+ /// `BIG_DECIMAL` ranges always become comparisons: an intermediate stage
matches `IN` values with `equals`, which
+ /// depends on the scale of the value.
+ private static RexExpression convertRanges(ColumnDataType dataType, RexNode
leftOperand, RangeSet rangeSet) {
+ Set<Range> ranges = rangeSet.asRanges();
+ if (dataType == ColumnDataType.BIG_DECIMAL) {
+ return convertRangesToOr(dataType, leftOperand, ranges);
+ }
+ List<Range> points = new ArrayList<>();
+ List<Range> others = new ArrayList<>();
+ splitPoints(ranges, points, others);
+ List<Range> complementPoints = new ArrayList<>();
+ List<Range> complementOthers = new ArrayList<>();
+ splitPoints(rangeSet.complement().asRanges(), complementPoints,
complementOthers);
+ boolean inForm = points.size() >= MIN_POINTS_FOR_IN;
+ boolean notInForm = complementPoints.size() >= MIN_POINTS_FOR_IN;
+ if (inForm && notInForm) {
+ // Use the form with fewer values. A range costs up to 2 values.
+ inForm = points.size() + 2 * others.size() <= complementPoints.size() +
2 * complementOthers.size();
+ notInForm = !inForm;
+ }
+ if (inForm) {
+ List<RexExpression> terms = new ArrayList<>(1 + others.size());
+ terms.add(new RexExpression.FunctionCall(ColumnDataType.BOOLEAN,
SqlKind.IN.name(),
+ toSearchFunctionOperands(leftOperand, points, dataType)));
+ for (Range range : others) {
+ terms.add(convertRangesToOr(dataType, leftOperand, Set.of(range)));
+ }
+ return combine(SqlKind.OR, terms);
+ }
+ if (notInForm) {
+ List<RexExpression> terms = new ArrayList<>(1 + complementOthers.size());
+ terms.add(new RexExpression.FunctionCall(ColumnDataType.BOOLEAN,
SqlKind.NOT_IN.name(),
+ toSearchFunctionOperands(leftOperand, complementPoints, dataType)));
+ for (Range range : complementOthers) {
+ // x is not in the range: x is in one of the (at most two) ranges of
the range's complement.
+ terms.add(convertRangesToOr(dataType, leftOperand,
ImmutableRangeSet.of(range).complement().asRanges()));
+ }
+ return combine(SqlKind.AND, terms);
+ }
+ return convertRangesToOr(dataType, leftOperand, ranges);
+ }
+
+ private static RexExpression combine(SqlKind kind, List<RexExpression>
terms) {
+ if (terms.size() == 1) {
+ return terms.get(0);
+ }
+ return new RexExpression.FunctionCall(ColumnDataType.BOOLEAN, kind.name(),
terms);
+ }
+
+ private static void splitPoints(Set<Range> ranges, List<Range> points,
List<Range> others) {
+ for (Range range : ranges) {
+ if (isPoint(range)) {
+ points.add(range);
+ } else {
+ others.add(range);
+ }
}
}
+ private static boolean isPoint(Range range) {
+ return range.hasLowerBound() && range.hasUpperBound() &&
range.lowerBoundType() == BoundType.CLOSED
+ && range.upperBoundType() == BoundType.CLOSED &&
range.lowerEndpoint().compareTo(range.upperEndpoint()) == 0;
+ }
+
private static RexExpression evaluateLiteralIn(RexLiteral leftOperand,
Set<Range> ranges, RexUnknownAs nullAs) {
// No need to do normal evaluation if the literal is a null literal and
nulls need to be included/excluded, so we
// can return early. Otherwise, continue with normal evaluation
@@ -512,8 +597,8 @@ public class RexExpressionUtils {
List.of(leftOperand, fromRexLiteralValue(dataType,
range.upperEndpoint())));
}
- /// Transforms a set of **point based** ranges into a list of expressions.
- private static List<RexExpression> toSearchFunctionOperands(RexNode
leftOperand, Set<Range> ranges,
+ /// Transforms a collection of **point based** ranges into a list of
expressions.
+ private static List<RexExpression> toSearchFunctionOperands(RexNode
leftOperand, Collection<Range> ranges,
ColumnDataType dataType) {
List<RexExpression> operands = new ArrayList<>(1 + ranges.size());
operands.add(fromRexNode(leftOperand));
diff --git
a/pinot-query-planner/src/test/java/org/apache/pinot/calcite/rex/SearchSealerTest.java
b/pinot-query-planner/src/test/java/org/apache/pinot/calcite/rex/SearchSealerTest.java
new file mode 100644
index 00000000000..ad48f7de5c7
--- /dev/null
+++
b/pinot-query-planner/src/test/java/org/apache/pinot/calcite/rex/SearchSealerTest.java
@@ -0,0 +1,315 @@
+/**
+ * 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.calcite.rex;
+
+import com.google.common.collect.ImmutableRangeSet;
+import com.google.common.collect.Range;
+import com.google.common.collect.RangeSet;
+import com.google.common.collect.TreeRangeSet;
+import java.math.BigDecimal;
+import java.util.ArrayList;
+import java.util.List;
+import org.apache.calcite.avatica.util.ByteString;
+import org.apache.calcite.plan.RelOptCluster;
+import org.apache.calcite.plan.RelOptPredicateList;
+import org.apache.calcite.plan.Strong;
+import org.apache.calcite.plan.hep.HepPlanner;
+import org.apache.calcite.plan.hep.HepProgram;
+import org.apache.calcite.rel.RelNode;
+import org.apache.calcite.rel.core.Filter;
+import org.apache.calcite.rel.logical.LogicalFilter;
+import org.apache.calcite.rel.logical.LogicalValues;
+import org.apache.calcite.rel.type.RelDataType;
+import org.apache.calcite.rex.RexBuilder;
+import org.apache.calcite.rex.RexCall;
+import org.apache.calcite.rex.RexLiteral;
+import org.apache.calcite.rex.RexNode;
+import org.apache.calcite.rex.RexSimplify;
+import org.apache.calcite.rex.RexUnknownAs;
+import org.apache.calcite.rex.RexUtil;
+import org.apache.calcite.rex.RexVisitorImpl;
+import org.apache.calcite.sql.SqlKind;
+import org.apache.calcite.sql.fun.SqlStdOperatorTable;
+import org.apache.calcite.sql.type.SqlTypeName;
+import org.apache.calcite.util.ImmutableBitSet;
+import org.apache.calcite.util.Sarg;
+import org.apache.calcite.util.TimestampString;
+import org.apache.pinot.query.planner.logical.RexExpression;
+import org.apache.pinot.query.planner.logical.RexExpressionUtils;
+import org.apache.pinot.query.type.TypeFactory;
+import org.testng.annotations.DataProvider;
+import org.testng.annotations.Test;
+
+import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertFalse;
+import static org.testng.Assert.assertNotEquals;
+import static org.testng.Assert.assertNotSame;
+import static org.testng.Assert.assertSame;
+import static org.testng.Assert.assertTrue;
+
+
+/// Unit tests for [SearchSealer] and [PinotSealedSearchOperator] on single
`SEARCH` calls. Planning tests are in
+/// `SealedInListPlanningTest`.
+public class SearchSealerTest {
+ private static final TypeFactory TYPE_FACTORY = new TypeFactory();
+ private static final RexBuilder REX_BUILDER = new RexBuilder(TYPE_FACTORY);
+
+ enum Shape {
+ POINTS, COMPLEMENTED_POINTS, RANGES, POINTS_AND_RANGES
+ }
+
+ @DataProvider(name = "searches")
+ public Object[][] searches() {
+ List<Object[]> cases = new ArrayList<>();
+ for (SqlTypeName type : new SqlTypeName[]{
+ SqlTypeName.INTEGER, SqlTypeName.BIGINT, SqlTypeName.DOUBLE,
SqlTypeName.DECIMAL, SqlTypeName.VARCHAR,
+ SqlTypeName.VARBINARY, SqlTypeName.TIMESTAMP, SqlTypeName.BOOLEAN
+ }) {
+ for (Shape shape : Shape.values()) {
+ if (type == SqlTypeName.BOOLEAN && shape != Shape.POINTS && shape !=
Shape.COMPLEMENTED_POINTS) {
+ continue;
+ }
+ for (RexUnknownAs nullAs : RexUnknownAs.values()) {
+ for (boolean nullable : new boolean[]{false, true}) {
+ cases.add(new Object[]{type, shape, nullAs, nullable});
+ }
+ }
+ }
+ }
+ return cases.toArray(new Object[0][]);
+ }
+
+ /// Sealing and un-sealing gives back the same `SEARCH`, and
`RexExpressionUtils` converts the sealed call, the
+ /// un-sealed call and the original call to the same expression.
+ @Test(dataProvider = "searches")
+ public void testRoundTrip(SqlTypeName typeName, Shape shape, RexUnknownAs
nullAs, boolean nullable) {
+ RexCall search = search(typeName, shape, nullAs, nullable);
+ SearchSealer sealer = new SearchSealer(1);
+ RexNode sealed = sealer.seal(REX_BUILDER, search);
+ assertTrue(sealed instanceof RexCall);
+ RexCall sealedCall = (RexCall) sealed;
+ assertTrue(sealedCall.getOperator() instanceof PinotSealedSearchOperator,
sealed.toString());
+ assertEquals(sealedCall.getKind(), SqlKind.OTHER_FUNCTION);
+ assertEquals(sealedCall.getOperands(),
List.of(search.getOperands().get(0)));
+ assertEquals(sealed.getType(), search.getType());
+ // The digest does not hold the values.
+ assertEquals(sealed.toString(), "$SEARCH#0($0)");
+
+ RexNode unsealed = SearchSealer.unsealCall(REX_BUILDER, sealedCall);
+ assertEquals(unsealed, search);
+ RexExpression expected = RexExpressionUtils.fromRexNode(search);
+ assertEquals(RexExpressionUtils.fromRexNode(unsealed), expected);
+ assertEquals(RexExpressionUtils.fromRexNode(sealed), expected);
+ }
+
+ /// The operator derives the same return type as `SEARCH` and gives the same
answers to [Strong].
+ @Test(dataProvider = "searches")
+ public void testSameTypeAndNullSemanticsAsSearch(SqlTypeName typeName, Shape
shape, RexUnknownAs nullAs,
+ boolean nullable) {
+ RexCall search = search(typeName, shape, nullAs, nullable);
+ RexCall sealed = (RexCall) new SearchSealer(1).seal(REX_BUILDER, search);
+ // Type derived from the operands, as RexBuilder#makeCall does when a rule
rebuilds the call.
+ RexNode rebuilt = REX_BUILDER.makeCall(sealed.getOperator(),
sealed.getOperands());
+ RexNode searchRebuilt = REX_BUILDER.makeCall(SqlStdOperatorTable.SEARCH,
search.getOperands());
+ assertEquals(rebuilt.getType(), searchRebuilt.getType());
+
+ ImmutableBitSet nullColumns = ImmutableBitSet.of(0);
+ assertEquals(Strong.isNull(sealed, nullColumns), Strong.isNull(search,
nullColumns));
+ assertEquals(Strong.isNotTrue(sealed, nullColumns),
Strong.isNotTrue(search, nullColumns));
+ assertEquals(Strong.isStrong(sealed), Strong.isStrong(search));
+ assertTrue(RexUtil.isDeterministic(sealed));
+ assertFalse(sealed.getOperator().isDynamicFunction());
+ assertTrue(sealed.getOperator().isSafeOperator());
+ }
+
+ @Test
+ public void testThreshold() {
+ RexCall search = search(SqlTypeName.INTEGER, Shape.POINTS,
RexUnknownAs.UNKNOWN, true);
+ int numRanges = search.getOperands().get(1).accept(new
RexVisitorImpl<Integer>(false) {
+ @Override
+ public Integer visitLiteral(RexLiteral literal) {
+ return literal.getValueAs(Sarg.class).rangeSet.asRanges().size();
+ }
+ });
+ assertSame(new SearchSealer(numRanges + 1).seal(REX_BUILDER, search),
search);
+ assertSame(new SearchSealer(0).seal(REX_BUILDER, search), search);
+ assertSame(new SearchSealer(-1).seal(REX_BUILDER, search), search);
+ assertNotSame(new SearchSealer(numRanges).seal(REX_BUILDER, search),
search);
+ }
+
+ /// One query shares one operator per distinct Sarg (so equal lists stay
equal), different Sargs get different
+ /// operators and names (so the digests differ), and operators of different
queries are never equal.
+ @Test
+ public void testOperatorIdentity() {
+ RexCall points = search(SqlTypeName.INTEGER, Shape.POINTS,
RexUnknownAs.UNKNOWN, true);
+ RexCall other = search(SqlTypeName.INTEGER, Shape.COMPLEMENTED_POINTS,
RexUnknownAs.UNKNOWN, true);
+ SearchSealer sealer = new SearchSealer(1);
+ RexCall sealed1 = (RexCall) sealer.seal(REX_BUILDER, points);
+ RexCall sealed2 = (RexCall) sealer.seal(REX_BUILDER, copy(points));
+ RexCall sealed3 = (RexCall) sealer.seal(REX_BUILDER, other);
+ assertSame(sealed1.getOperator(), sealed2.getOperator());
+ assertEquals(sealed1, sealed2);
+ assertNotEquals(sealed1.getOperator(), sealed3.getOperator());
+ assertNotEquals(sealed1.toString(), sealed3.toString());
+
+ RexCall otherQuery = (RexCall) new SearchSealer(1).seal(REX_BUILDER,
points);
+ assertEquals(otherQuery.getOperator().getName(),
sealed1.getOperator().getName());
+ assertNotEquals(otherQuery.getOperator(), sealed1.getOperator());
+ assertNotEquals(otherQuery, sealed1);
+ }
+
+ /// When the operand becomes a literal (for example after a filter is pushed
through a project of constants),
+ /// simplification keeps the call and `RexExpressionUtils` evaluates it like
`SEARCH(literal, sarg)`.
+ @Test
+ public void testLiteralOperand() {
+ RexCall search = search(SqlTypeName.INTEGER, Shape.POINTS,
RexUnknownAs.UNKNOWN, true);
+ RexCall sealed = (RexCall) new SearchSealer(1).seal(REX_BUILDER, search);
+ for (RexLiteral literal : new RexLiteral[]{
+ REX_BUILDER.makeExactLiteral(BigDecimal.valueOf(3)),
REX_BUILDER.makeExactLiteral(BigDecimal.valueOf(4)),
+ REX_BUILDER.makeNullLiteral(search.getOperands().get(0).getType())
+ }) {
+ RexCall withLiteral = sealed.clone(sealed.getType(), List.of(literal));
+ RexNode simplified = new RexSimplify(REX_BUILDER,
RelOptPredicateList.EMPTY, RexUtil.EXECUTOR)
+ .simplifyUnknownAs(withLiteral, RexUnknownAs.UNKNOWN);
+ RexCall searchWithLiteral = search.clone(search.getType(),
List.of(literal, search.getOperands().get(1)));
+ if (literal.isNull()) {
+ // Like SEARCH with NULL AS UNKNOWN, the call is null when its operand
is null.
+ assertTrue(RexUtil.isNullLiteral(simplified, true),
simplified.toString());
+ } else {
+ assertSame(simplified, withLiteral);
+ }
+ assertEquals(RexExpressionUtils.fromRexNode(withLiteral),
RexExpressionUtils.fromRexNode(searchWithLiteral));
+ }
+ }
+
+ /// Many comparisons of one operand fold into a Sarg, which is then sealed.
Comparisons of different operands cannot
+ /// fold into one Sarg, so a filter with many of them is left as it is.
+ @Test
+ public void testFoldsOnlyComparisonsOfOneOperand() {
+ RelOptCluster cluster = RelOptCluster.create(new
HepPlanner(HepProgram.builder().build()), REX_BUILDER);
+ RelDataType intType = TYPE_FACTORY.createSqlType(SqlTypeName.INTEGER);
+ RelNode values = LogicalValues.createEmpty(cluster,
TYPE_FACTORY.builder().add("x", intType).build());
+ RexNode x = REX_BUILDER.makeInputRef(values, 0);
+ List<RexNode> sameOperand = new ArrayList<>();
+ List<RexNode> differentOperands = new ArrayList<>();
+ for (int i = 1; i <= 25; i++) {
+ RexNode literal = REX_BUILDER.makeExactLiteral(BigDecimal.valueOf(i),
intType);
+ sameOperand.add(REX_BUILDER.makeCall(SqlStdOperatorTable.NOT_EQUALS, x,
literal));
+ differentOperands.add(REX_BUILDER.makeCall(SqlStdOperatorTable.EQUALS,
+ REX_BUILDER.makeCall(SqlStdOperatorTable.PLUS, x, literal),
literal));
+ }
+ RelNode list = LogicalFilter.create(values,
REX_BUILDER.makeCall(SqlStdOperatorTable.AND, sameOperand));
+ String sealedCondition = ((Filter) new
SearchSealer(20).seal(list)).getCondition().toString();
+
assertTrue(sealedCondition.startsWith(PinotSealedSearchOperator.NAME_PREFIX),
sealedCondition);
+
+ // Simplification would remove the repeated term, so the filter only stays
the same if it is not simplified.
+ differentOperands.add(differentOperands.get(0));
+ RelNode wide = LogicalFilter.create(values,
REX_BUILDER.makeCall(SqlStdOperatorTable.AND, differentOperands));
+ assertSame(new SearchSealer(20).seal(wide), wide);
+ }
+
+ private static RexCall copy(RexCall search) {
+ RexLiteral literal = (RexLiteral) search.getOperands().get(1);
+ return (RexCall) REX_BUILDER.makeCall(SqlStdOperatorTable.SEARCH,
search.getOperands().get(0),
+ REX_BUILDER.makeSearchArgumentLiteral(literal.getValueAs(Sarg.class),
literal.getType()));
+ }
+
+ @SuppressWarnings({"rawtypes", "unchecked"})
+ static RexCall search(SqlTypeName typeName, Shape shape, RexUnknownAs
nullAs, boolean nullable) {
+ List<RexLiteral> literals = literals(typeName);
+ RelDataType literalType = literals.get(0).getType();
+ RelDataType operandType =
+
TYPE_FACTORY.createTypeWithNullability(TYPE_FACTORY.createSqlType(typeName),
nullable);
+ RexNode operand = REX_BUILDER.makeInputRef(operandType, 0);
+ RangeSet rangeSet = TreeRangeSet.create();
+ List<Comparable> values = new ArrayList<>();
+ for (RexLiteral literal : literals) {
+ values.add(literal.getValueAs(Comparable.class));
+ }
+ switch (shape) {
+ case POINTS:
+ values.forEach(v -> rangeSet.add(Range.singleton(v)));
+ break;
+ case COMPLEMENTED_POINTS:
+ values.forEach(v -> rangeSet.add(Range.singleton(v)));
+ break;
+ case RANGES:
+ for (int i = 0; i + 1 < values.size(); i += 2) {
+ rangeSet.add(Range.open(values.get(i), values.get(i + 1)));
+ }
+ break;
+ case POINTS_AND_RANGES:
+ rangeSet.add(Range.singleton(values.get(0)));
+ rangeSet.add(Range.closedOpen(values.get(1), values.get(2)));
+ rangeSet.add(Range.singleton(values.get(3)));
+ rangeSet.add(Range.greaterThan(values.get(values.size() - 1)));
+ break;
+ default:
+ throw new IllegalStateException();
+ }
+ RangeSet finalRangeSet = shape == Shape.COMPLEMENTED_POINTS ?
rangeSet.complement() : rangeSet;
+ Sarg sarg = Sarg.of(nullAs, ImmutableRangeSet.copyOf(finalRangeSet));
+ return (RexCall) REX_BUILDER.makeCall(SqlStdOperatorTable.SEARCH, operand,
+ REX_BUILDER.makeSearchArgumentLiteral(sarg, literalType));
+ }
+
+ private static List<RexLiteral> literals(SqlTypeName typeName) {
+ List<RexLiteral> literals = new ArrayList<>();
+ switch (typeName) {
+ case BOOLEAN:
+ literals.add(REX_BUILDER.makeLiteral(false));
+ literals.add(REX_BUILDER.makeLiteral(true));
+ return literals;
+ default:
+ break;
+ }
+ for (int i = 1; i <= 6; i++) {
+ switch (typeName) {
+ case INTEGER:
+ literals.add(REX_BUILDER.makeExactLiteral(BigDecimal.valueOf(i * 3),
+ TYPE_FACTORY.createSqlType(SqlTypeName.INTEGER)));
+ break;
+ case BIGINT:
+
literals.add(REX_BUILDER.makeExactLiteral(BigDecimal.valueOf(10_000_000_000L +
i),
+ TYPE_FACTORY.createSqlType(SqlTypeName.BIGINT)));
+ break;
+ case DOUBLE:
+ literals.add(REX_BUILDER.makeApproxLiteral(BigDecimal.valueOf(i +
0.5),
+ TYPE_FACTORY.createSqlType(SqlTypeName.DOUBLE)));
+ break;
+ case DECIMAL:
+ literals.add(REX_BUILDER.makeExactLiteral(new BigDecimal(i +
".25")));
+ break;
+ case VARCHAR:
+ literals.add(REX_BUILDER.makeLiteral("value-" + i));
+ break;
+ case VARBINARY:
+ literals.add(REX_BUILDER.makeBinaryLiteral(new ByteString(new
byte[]{(byte) i, 0x7f})));
+ break;
+ case TIMESTAMP:
+ long millis = 1_700_000_000_000L + i * 60_000L;
+
literals.add(REX_BUILDER.makeTimestampLiteral(TimestampString.fromMillisSinceEpoch(millis),
1));
+ break;
+ default:
+ throw new IllegalStateException();
+ }
+ }
+ return literals;
+ }
+}
diff --git
a/pinot-query-planner/src/test/java/org/apache/pinot/query/SealedInListPlanningTest.java
b/pinot-query-planner/src/test/java/org/apache/pinot/query/SealedInListPlanningTest.java
new file mode 100644
index 00000000000..756bfe3aa64
--- /dev/null
+++
b/pinot-query-planner/src/test/java/org/apache/pinot/query/SealedInListPlanningTest.java
@@ -0,0 +1,555 @@
+/**
+ * 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.query;
+
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.stream.Collectors;
+import java.util.stream.IntStream;
+import javax.annotation.Nullable;
+import org.apache.calcite.plan.RelOptRule;
+import org.apache.calcite.plan.RelOptRuleCall;
+import org.apache.calcite.rel.RelNode;
+import org.apache.calcite.rex.RexCall;
+import org.apache.calcite.rex.RexLiteral;
+import org.apache.calcite.rex.RexNode;
+import org.apache.calcite.rex.RexShuttle;
+import org.apache.calcite.sql.SqlBasicCall;
+import org.apache.calcite.sql.SqlCall;
+import org.apache.calcite.sql.SqlKind;
+import org.apache.calcite.sql.fun.SqlStdOperatorTable;
+import org.apache.calcite.sql.util.SqlBasicVisitor;
+import org.apache.calcite.util.Sarg;
+import org.apache.pinot.common.config.provider.TableCache;
+import org.apache.pinot.core.routing.MockRoutingManagerFactory;
+import org.apache.pinot.core.routing.RoutingManager;
+import org.apache.pinot.query.planner.physical.DispatchablePlanFragment;
+import org.apache.pinot.query.planner.physical.DispatchableSubPlan;
+import org.apache.pinot.query.planner.rules.DefaultRuleSetCustomizer;
+import org.apache.pinot.query.planner.rules.PinotRuleSet;
+import org.apache.pinot.query.planner.serde.PlanNodeSerializer;
+import org.apache.pinot.query.planner.spi.Phase;
+import org.apache.pinot.query.planner.spi.RuleSetCustomizer;
+import org.apache.pinot.query.routing.WorkerManager;
+import org.apache.pinot.spi.data.Schema;
+import org.apache.pinot.spi.utils.CommonConstants;
+import org.apache.pinot.sql.parsers.CalciteSqlParser;
+import org.apache.pinot.sql.parsers.SqlNodeAndOptions;
+import org.testng.annotations.BeforeClass;
+import org.testng.annotations.DataProvider;
+import org.testng.annotations.Test;
+
+import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertFalse;
+import static org.testng.Assert.assertTrue;
+import static org.testng.Assert.expectThrows;
+
+
+/// Planning with sealed IN lists (see `SearchSealer`) must give the same plan
as planning without sealing, must not
+/// let any planner rule see a large Sarg, and must be fast in the positions
where Calcite is super-linear.
+public class SealedInListPlanningTest extends QueryEnvironmentTestBase {
+ private static final int NUM_VALUES = 50;
+ /// Threshold of the guarded planner, and the smallest Sarg or comparison
list that the guard rejects.
+ private static final int GUARD_THRESHOLD = 20;
+ private static final String INTS = IntStream.range(0, NUM_VALUES).mapToObj(i
-> Integer.toString(i * 7 + 3))
+ .collect(Collectors.joining(", "));
+ private static final String STRINGS = IntStream.range(0,
NUM_VALUES).mapToObj(i -> "'v" + i * 7 + "'")
+ .collect(Collectors.joining(", "));
+ private static final String OTHER_INTS = IntStream.range(0,
NUM_VALUES).mapToObj(i -> Integer.toString(i * 5 + 3))
+ .collect(Collectors.joining(", "));
+ private static final String TIMESTAMPS = IntStream.range(0, 25)
+ .mapToObj(i -> String.format("TIMESTAMP '2024-01-01 00:00:%02d.%03d'",
i, i)).collect(Collectors.joining(", "));
+ private static final String BOOLEANS = IntStream.range(0, 24).mapToObj(i ->
i % 2 == 0 ? "TRUE" : "FALSE")
+ .collect(Collectors.joining(", "));
+ private static final String OR_CHAIN = IntStream.range(0,
NUM_VALUES).mapToObj(i -> "col3 = " + (i * 7 + 3))
+ .collect(Collectors.joining(" OR "));
+ /// Comparisons of many different operands, which cannot fold into one Sarg.
+ private static final String WIDE_AND = IntStream.range(1, 25).mapToObj(i ->
"col3 + " + i + " = " + i * 7)
+ .collect(Collectors.joining(" AND "));
+ private static final String WIDE_OR = IntStream.range(1, 25).mapToObj(i ->
"col6 + " + i + " = " + i * 7)
+ .collect(Collectors.joining(" OR "));
+
+ /// The same tables with column-based null handling, so that their columns
are nullable.
+ private QueryEnvironment _nullableQueryEnvironment;
+ /// The nullable tables, planned as the broker plans queries without null
handling.
+ private QueryEnvironment _withoutNullHandling;
+ private QueryEnvironment _guardedWithoutNullHandling;
+
+ @BeforeClass
+ public void setUpNullableTables() {
+ Map<String, Schema> schemas = new HashMap<>();
+ for (String table : List.of("a", "b", "c")) {
+ String tableName = TABLE_SCHEMAS.containsKey(table + "_REALTIME") ?
table + "_REALTIME" : table + "_OFFLINE";
+ schemas.put(tableName,
getSchemaBuilder(table).setEnableColumnBasedNullHandling(true).build());
+ }
+ _nullableQueryEnvironment =
+ getQueryEnvironment(3, 1, 2, schemas, SERVER1_SEGMENTS,
SERVER2_SEGMENTS, PARTITIONED_SEGMENTS_MAP);
+ Map<String, Schema> allSchemas = new HashMap<>(TABLE_SCHEMAS);
+ allSchemas.putAll(schemas);
+ _withoutNullHandling = newQueryEnvironment(allSchemas, false);
+ _guardedWithoutNullHandling = newQueryEnvironment(allSchemas, true);
+ }
+
+ @DataProvider(name = "samePlanQueries")
+ public Object[][] samePlanQueries() {
+ List<String> queries = List.of(
+ "SELECT col1, col3 FROM a WHERE col3 IN (" + INTS + ")",
+ "SELECT * FROM a WHERE col3 IN (" + INTS + ")",
+ "SELECT col1 FROM a WHERE col3 NOT IN (" + INTS + ")",
+ "SELECT col1 FROM a WHERE col3 IN (" + INTS + ", NULL)",
+ "SELECT col1 FROM a WHERE col3 NOT IN (" + INTS + ", NULL)",
+ "SELECT col1 FROM a WHERE col3 IN (col6, " + INTS + ")",
+ "SELECT col1 FROM a WHERE col3 NOT IN (col6, " + INTS + ")",
+ "SELECT col1 FROM a WHERE col1 IN (" + STRINGS + ")",
+ "SELECT col1 FROM a WHERE UPPER(col1) IN (" + STRINGS + ")",
+ "SELECT col1 FROM a WHERE col7 IN (" + INTS + ")",
+ "SELECT col1 FROM a WHERE col4 IN (" + INTS + ")",
+ "SELECT col1 FROM a WHERE col3 IN (" + STRINGS.replace("'v", "'") +
")",
+ "SELECT col1 FROM a WHERE NOT (col3 IN (" + INTS + ") AND col1 = 'x')",
+ "SELECT col1 FROM a WHERE " + OR_CHAIN,
+ // All fields stay, so field trimming does not fold the chain before
sealing.
+ "SELECT * FROM a WHERE " + OR_CHAIN,
+ "SELECT * FROM a WHERE " + OR_CHAIN.replaceAll("col3 = (\\d+)", "$1 =
col3"),
+ "SELECT SUM(CASE WHEN col3 IN (" + INTS + ") THEN 1 ELSE 0 END) FROM
a",
+ "SELECT SUM(CASE WHEN col3 NOT IN (" + INTS + ") THEN 1 ELSE 0 END)
FROM a",
+ "SELECT COUNT(*) FILTER (WHERE col3 IN (" + INTS + ")), COUNT(*) FROM
a",
+ "SELECT col3 IN (" + INTS + ") FROM a",
+ "SELECT CASE WHEN col3 IN (" + INTS + ") THEN 'x' ELSE 'y' END,
COUNT(*) FROM a "
+ + "GROUP BY CASE WHEN col3 IN (" + INTS + ") THEN 'x' ELSE 'y'
END",
+ "SELECT col1, COUNT(*) FROM a GROUP BY col1 HAVING COUNT(*) IN (" +
INTS + ")",
+ "SELECT col1, ROW_NUMBER() OVER (PARTITION BY CASE WHEN col3 IN (" +
INTS + ") THEN 1 ELSE 0 END "
+ + "ORDER BY col6) FROM a",
+ "SELECT a.col1 FROM a JOIN b ON a.col1 = b.col1 AND b.col3 IN (" +
INTS + ")",
+ "SELECT a.col1 FROM a LEFT JOIN b ON a.col1 = b.col1 AND a.col3 IN ("
+ INTS + ")",
+ // The IN list rejects nulls, so the LEFT JOIN becomes an INNER JOIN.
+ "SELECT a.col1 FROM a LEFT JOIN b ON a.col1 = b.col1 WHERE b.col3 IN
(" + INTS + ")",
+ "SELECT a.col1 FROM a LEFT JOIN b ON a.col1 = b.col1 WHERE b.col3 NOT
IN (" + INTS + ")",
+ // The IN list on the join key is copied to the other side by
JoinPushTransitivePredicates.
+ "SELECT a.col1, b.col2 FROM a JOIN b ON a.col3 = b.col3 WHERE a.col3
IN (" + INTS + ")",
+ "SELECT a.col1, b.col2 FROM a LEFT JOIN b ON a.col3 = b.col3 WHERE
a.col3 NOT IN (" + INTS + ")",
+ "WITH t AS (SELECT col1, col3 FROM a WHERE col3 IN (" + INTS + ")) "
+ + "SELECT t1.col1 FROM t t1 JOIN t t2 ON t1.col1 = t2.col1",
+ "SELECT col1 FROM a WHERE col3 IN (" + INTS + ") UNION ALL SELECT col1
FROM b WHERE col3 IN (" + INTS + ")",
+ "SELECT col1 FROM a WHERE col3 IN (SELECT col3 FROM b WHERE col6 IN ("
+ INTS + "))",
+ "SELECT col1 FROM a WHERE col3 IN (" + INTS + ") AND col6 IN (" + INTS
+ ")",
+ "SELECT col2, COUNT(*) FROM a WHERE col3 IN (" + INTS + ") GROUP BY
col2 ORDER BY COUNT(*) DESC LIMIT 5",
+ "SELECT a.col2, b.col2, COUNT(*) FILTER (WHERE a.col1 = 'x') FROM a
LEFT JOIN b ON a.col1 = b.col1 "
+ + "LEFT JOIN c ON a.col2 = c.col2 WHERE a.col3 IN (" + INTS + ")
AND a.ts > 10 "
+ + "GROUP BY a.col2, b.col2 ORDER BY 3 DESC LIMIT 50",
+ // Correlated sub-queries (the list is sealed after decorrelation).
+ "SELECT col1 FROM a WHERE EXISTS (SELECT 1 FROM b WHERE b.col1 =
a.col1 AND b.col3 IN (" + INTS + "))",
+ "SELECT a.col1, (SELECT MAX(b.col6) FROM b WHERE b.col1 = a.col1 AND
b.col3 IN (" + INTS + ")) FROM a",
+ "SELECT col1 FROM a WHERE col3 NOT IN (SELECT col3 FROM b WHERE col6
IN (" + INTS + "))",
+ // Set operations, sort, partition column, hints.
+ "SELECT col1 FROM a WHERE col3 IN (" + INTS + ") INTERSECT SELECT col1
FROM b WHERE col3 NOT IN (" + INTS
+ + ")",
+ "SELECT col1 FROM a WHERE col3 IN (" + INTS + ") EXCEPT SELECT col1
FROM b WHERE col6 IN (" + INTS + ")",
+ "SELECT * FROM a WHERE col3 IN (" + INTS + ") ORDER BY col1 LIMIT 5
OFFSET 2",
+ "SELECT col1 FROM a WHERE col2 IN (" + STRINGS + ")",
+ "SELECT /*+ aggOptions(is_partitioned_by_group_by_keys='true') */
col2, COUNT(*) FROM a "
+ + "WHERE col3 IN (" + INTS + ") GROUP BY col2",
+ "SELECT /*+ joinOptions(join_strategy = 'lookup') */ a.col1, b.col2
FROM a JOIN b ON a.col1 = b.col1 "
+ + "WHERE a.col3 IN (" + INTS + ")",
+ "SELECT /*+ joinOptions(join_strategy = 'hash') */ col1 FROM a WHERE
col3 IN (SELECT col3 FROM b "
+ + "WHERE col6 IN (" + INTS + "))",
+ "SELECT a.col1 FROM a /*+ tableOptions(partition_function='hashcode',
partition_key='col2', "
+ + "partition_size='4') */ JOIN b /*+
tableOptions(partition_function='hashcode', partition_key='col1', "
+ + "partition_size='4') */ ON a.col2 = b.col1 WHERE a.col3 IN (" +
INTS + ")",
+ "SET useSpools=true; WITH t AS (SELECT col1, col3 FROM a WHERE col3 IN
(" + INTS + ")) "
+ + "SELECT t1.col1 FROM t t1 JOIN t t2 ON t1.col1 = t2.col1",
+ "SET useLiteMode=true; SELECT col1 FROM a WHERE col3 IN (" + INTS + ")
LIMIT 10",
+ // Other predicates on the same column fold into the list before it is
sealed, as without sealing.
+ "SELECT col1 FROM a WHERE col3 IN (" + INTS + ") AND col3 IN (" +
OTHER_INTS + ")",
+ "SELECT col1 FROM a WHERE col3 IN (" + INTS + ") OR col3 IN (" +
OTHER_INTS + ")",
+ "SELECT col1 FROM a WHERE col3 IN (" + INTS + ") AND col3 = 10",
+ "SELECT col1 FROM a WHERE col3 IN (" + INTS + ") AND col3 = 3",
+ "SELECT col1 FROM a WHERE col3 IN (" + INTS + ") OR col3 > 1000000",
+ "SELECT col1 FROM a WHERE col3 NOT IN (" + INTS + ") AND col3 > 0",
+ "SELECT col1 FROM a WHERE col3 IN (" + INTS + ") OR col3 BETWEEN 1000
AND 2000",
+ // NOT over a sealed call, which un-sealing folds into the negated
Sarg.
+ "SELECT NOT (col3 IN (" + INTS + ")) FROM a",
+ "SELECT COUNT(CASE WHEN col3 IN (" + INTS + ") THEN NULL ELSE 1 END)
FROM a",
+ // Joins keep their distribution after un-sealing, so the exchange
above stays pre-partitioned.
+ "SELECT a.col1, b.col1 FROM a JOIN b ON a.col1 = b.col1 JOIN c ON
a.col1 = c.col1 WHERE a.col3 IN (" + INTS
+ + ")",
+ "SELECT a.col1, COUNT(*) FROM a JOIN b ON a.col1 = b.col1 WHERE b.col3
IN (" + INTS + ") GROUP BY a.col1",
+ // Other types.
+ "SELECT col1 FROM a WHERE ts_timestamp IN (" + TIMESTAMPS + ")",
+ "SELECT col1 FROM a WHERE col5 IN (" + BOOLEANS + ")",
+ // Many comparisons of different operands, without an IN list.
+ "SELECT col1 FROM a WHERE " + WIDE_AND,
+ "SELECT col1 FROM a WHERE " + WIDE_OR,
+ "SELECT a.col1 FROM a JOIN b ON a.col1 = b.col1 AND (" +
WIDE_OR.replace("col6", "b.col6") + ")",
+ "SELECT col1 FROM a WHERE col3 IN (" + INTS + ") AND " + WIDE_AND
+ );
+ List<Object[]> cases = new ArrayList<>();
+ for (String query : queries) {
+ cases.add(new Object[]{query});
+ cases.add(new Object[]{"SET usePhysicalOptimizer=true; " + query});
+ }
+ return cases.toArray(new Object[0][]);
+ }
+
+ @Test(dataProvider = "samePlanQueries")
+ public void testSamePlanAsWithoutSealing(String query) {
+ String unsealed = "SET sealedInListThreshold=0; " + query;
+ String explained = explain(query);
+ assertFalse(explained.contains("$SEARCH#"), explained);
+ assertEquals(explained, explain(unsealed));
+ // The physical optimizer picks a random server for some single-worker
stages, so hosts are not compared.
+ assertEquals(implementationPlan(query), implementationPlan(unsealed));
+ assertEquals(serializedStages(query), serializedStages(unsealed));
+ }
+
+ /// Null checks on the column of a large list must give the same plan as
without sealing. This planner uses Calcite's
+ /// `IS NULL` and `IS NOT NULL`, as the broker does with null handling.
Calcite folds `x IS NOT NULL` into the Sarg of
+ /// `x NOT IN (...)` as `NULL AS FALSE`, and the lists are sealed after
this, so the plans stay the same.
+ ///
+ /// Not in this list: a null check in another clause than the list (for
example the list in `JOIN ... ON` and
+ /// `x IS NOT NULL` in `WHERE`) only meets the sealed list during
optimization. Calcite then drops the null check as
+ /// redundant, as it does next to `x < 5`. This is correct, because with
null handling the servers apply SQL null
+ /// semantics. [#testNullChecksWithoutNullHandling] covers these queries
without null handling.
+ @DataProvider(name = "nullableColumnQueries")
+ public Object[][] nullableColumnQueries() {
+ List<String> queries = List.of(
+ "SELECT col1 FROM a WHERE col3 NOT IN (" + INTS + ") AND col3 IS NOT
NULL",
+ "SELECT col1 FROM a WHERE col3 IN (" + INTS + ") AND col3 IS NOT NULL",
+ "SELECT col1 FROM a WHERE NOT (col3 IN (" + INTS + ") OR col3 IS
NULL)",
+ "SELECT col1 FROM a WHERE col3 IN (" + INTS + ") OR col3 IS NULL",
+ "SELECT col1 FROM a WHERE col3 NOT IN (" + INTS + ") OR col3 IN (" +
INTS + ")",
+ "SELECT col1 FROM a WHERE col3 NOT IN (" + INTS + ", NULL)",
+ "SELECT col1 FROM a WHERE col3 NOT IN (" + INTS + ") AND col6 IS NOT
NULL",
+ "SELECT col3 IN (" + INTS + ") OR col3 IS NULL, col3 NOT IN (" + INTS
+ ") FROM a",
+ "SELECT a.col1 FROM a LEFT JOIN b ON a.col1 = b.col1 WHERE b.col3 IN
(" + INTS + ")",
+ "SELECT a.col1 FROM a LEFT JOIN b ON a.col1 = b.col1 WHERE b.col3 NOT
IN (" + INTS + ") OR b.col3 IS NULL",
+ "SELECT COUNT(*) FILTER (WHERE col3 NOT IN (" + INTS + ") AND col3 IS
NOT NULL) FROM a"
+ );
+ List<Object[]> cases = new ArrayList<>();
+ for (String query : queries) {
+ for (String options : List.of("", "SET enableNullHandling=true; ", "SET
usePhysicalOptimizer=true; ")) {
+ cases.add(new Object[]{options + query});
+ }
+ }
+ return cases.toArray(new Object[0][]);
+ }
+
+ /// Null checks with the planner set up as the broker does it for queries
without null handling. The planner then uses
+ /// Pinot's `IS NULL` and `IS NOT NULL` operators, which Calcite cannot see
into. So a null check stays in the plan,
+ /// wherever it is, and the plan is the same as without sealing.
+ @DataProvider(name = "nullChecksWithoutNullHandling")
+ public Object[][] nullChecksWithoutNullHandling() {
+ List<String> queries = List.of(
+ // The null check is in another clause than the list.
+ "SELECT a.col1 FROM a JOIN b ON a.col1 = b.col1 AND b.col3 NOT IN (" +
INTS + ") WHERE b.col3 IS NOT NULL",
+ "SELECT a.col1 FROM a LEFT JOIN b ON a.col1 = b.col1 AND b.col3 NOT IN
(" + INTS + ") "
+ + "WHERE b.col3 IS NOT NULL",
+ "SELECT col1 FROM (SELECT col1, col3 FROM a WHERE col3 NOT IN (" +
INTS + ")) t WHERE t.col3 IS NOT NULL",
+ "SELECT col1 FROM (SELECT col1, col3 FROM a WHERE col3 IN (" + INTS +
")) t WHERE t.col3 IS NOT NULL",
+ "WITH t AS (SELECT col1, col3 FROM a WHERE col3 NOT IN (" + INTS + "))
SELECT col1 FROM t "
+ + "WHERE col3 IS NOT NULL",
+ // The null check is on a column that the column of the list is joined
on.
+ "SELECT a.col1 FROM a JOIN b ON a.col3 = b.col6 WHERE a.col3 NOT IN ("
+ INTS + ") AND b.col6 IS NOT NULL",
+ // The null check is in the same condition as the list.
+ "SELECT col1 FROM a WHERE col3 NOT IN (" + INTS + ") AND col3 IS NOT
NULL",
+ "SELECT col1 FROM a WHERE col3 IN (" + INTS + ") OR col3 IS NULL",
+ // NOT IN over a sub-query adds null checks of its own.
+ "SELECT col1 FROM (SELECT col1, col3 FROM a WHERE col3 NOT IN (" +
INTS + ")) t "
+ + "WHERE t.col3 NOT IN (SELECT col6 FROM b)"
+ );
+ List<Object[]> cases = new ArrayList<>();
+ for (String query : queries) {
+ cases.add(new Object[]{query});
+ cases.add(new Object[]{"SET usePhysicalOptimizer=true; " + query});
+ }
+ return cases.toArray(new Object[0][]);
+ }
+
+ @Test(dataProvider = "nullChecksWithoutNullHandling")
+ public void testNullChecksWithoutNullHandling(String query) {
+ String unsealed = "SET sealedInListThreshold=0; " + query;
+ String explained = explain(_withoutNullHandling, query);
+ if (query.contains(" IS NOT NULL") || query.contains(" IS NULL")) {
+ assertTrue(explained.contains(" NULL($"), explained);
+ }
+ assertEquals(explained, explain(_withoutNullHandling, unsealed));
+ assertEquals(serializedStages(_withoutNullHandling, query),
serializedStages(_withoutNullHandling, unsealed));
+ // The list is sealed: no planner rule sees it.
+ try (QueryEnvironment.CompiledQuery compiledQuery =
_guardedWithoutNullHandling.compile(query)) {
+ compiledQuery.planQuery(1);
+ }
+ }
+
+ @Test(dataProvider = "nullableColumnQueries")
+ public void testSamePlanAsWithoutSealingForNullableColumns(String query) {
+ String unsealed = "SET sealedInListThreshold=0; " + query;
+ assertEquals(explain(_nullableQueryEnvironment, query),
explain(_nullableQueryEnvironment, unsealed));
+ assertEquals(serializedStages(_nullableQueryEnvironment, query),
+ serializedStages(_nullableQueryEnvironment, unsealed));
+ }
+
+ /// An IN list and a range on the same column fold into one Sarg with points
and ranges, as without sealing. The
+ /// leaf gets one `IN` and the range, not one range per value.
+ @Test
+ public void testInListOrRangeShipsInAndRange() {
+ for (String options : List.of("", "SET usePhysicalOptimizer=true; ")) {
+ String stages = String.join("\n",
+ serializedStages(options + "SELECT col1 FROM a WHERE col3 IN (" +
INTS + ") OR col3 > 1000000"));
+ assertTrue(stages.contains("functionName: \"IN\""), stages);
+ assertTrue(stages.contains("functionName: \"GREATER_THAN\""), stages);
+ assertFalse(stages.contains("LESS_THAN_OR_EQUAL"), stages);
+ }
+ }
+
+ /// Rules can put a literal in place of the sealed operand, or combine the
sealed call with a contradicting
+ /// predicate. Planning must still work.
+ @Test
+ public void testSealedCallWithLiteralOperand() {
+ for (String query : List.of(
+ "SELECT x FROM (SELECT 5 AS x, col1 FROM a) WHERE x IN (" + INTS + ")",
+ "SELECT x FROM (SELECT 10 AS x, col1 FROM a) WHERE x IN (" + INTS +
")",
+ "SELECT col1 FROM a WHERE col3 = 5 AND col3 IN (" + INTS + ")",
+ "SELECT col1 FROM a WHERE col3 = 10 AND col3 IN (" + INTS + ")",
+ "SELECT col1 FROM a WHERE col3 IS NULL AND col3 IN (" + INTS + ")")) {
+ for (String prefix : List.of("", "SET usePhysicalOptimizer=true; ")) {
+ try (QueryEnvironment.CompiledQuery compiledQuery =
_queryEnvironment.compile(prefix + query)) {
+
assertFalse(compiledQuery.planQuery(1).getQueryPlan().getQueryStages().isEmpty());
+ }
+ }
+ }
+ }
+
+ /// The operators that the SQL tree had before sealing are put back after
the conversion to relational algebra.
+ @Test
+ public void testSqlTreeIsRestored() {
+ String query = "SELECT SUM(CASE WHEN col3 IN (" + INTS + ") THEN 1 ELSE 0
END) FROM a WHERE col3 NOT IN ("
+ + INTS + ")";
+ SqlNodeAndOptions sqlNodeAndOptions =
CalciteSqlParser.compileToSqlNodeAndOptions(query);
+ try (QueryEnvironment.CompiledQuery compiledQuery =
_queryEnvironment.compile(query, sqlNodeAndOptions)) {
+ List<SqlCall> inCalls = new ArrayList<>();
+ sqlNodeAndOptions.getSqlNode().accept(new SqlBasicVisitor<Void>() {
+ @Override
+ public Void visit(SqlCall call) {
+ if (call.getKind() == SqlKind.IN || call.getKind() ==
SqlKind.NOT_IN) {
+ inCalls.add(call);
+ }
+ return super.visit(call);
+ }
+ });
+ assertEquals(inCalls.size(), 2);
+ for (SqlCall call : inCalls) {
+ assertTrue(call instanceof SqlBasicCall);
+ assertTrue(call.getOperator() == SqlStdOperatorTable.IN ||
call.getOperator() == SqlStdOperatorTable.NOT_IN,
+ call.getOperator().getClass().getName());
+ }
+ }
+ }
+
+ /// Without sealing, these positions plan in cubic time (minutes for a few
thousand values).
+ @Test(timeOut = 60_000)
+ public void testLargeListsOutsideWhereAreFast() {
+ String values = IntStream.range(0,
5_000).mapToObj(Integer::toString).collect(Collectors.joining(", "));
+ for (String query : List.of(
+ "SELECT SUM(CASE WHEN col3 IN (" + values + ") THEN 1 ELSE 0 END) FROM
a",
+ "SELECT COUNT(*) FILTER (WHERE col3 NOT IN (" + values + ")) FROM a",
+ "SELECT col3 IN (" + values + ") FROM a",
+ "SELECT a.col1 FROM a JOIN b ON a.col1 = b.col1 AND b.col3 IN (" +
values + ")",
+ "SELECT col1, ROW_NUMBER() OVER (PARTITION BY CASE WHEN col3 IN (" +
values + ") THEN 1 ELSE 0 END) FROM a")) {
+ try (QueryEnvironment.CompiledQuery compiledQuery =
_queryEnvironment.compile(query)) {
+
assertFalse(compiledQuery.planQuery(1).getQueryPlan().getQueryStages().isEmpty());
+ }
+ }
+ }
+
+ /// A rule placed in every phase checks that no planner rule is offered a
plan with a large `SEARCH`, or with a
+ /// large `AND`/`OR` of comparisons that Calcite could fold into one.
+ @Test(dataProvider = "samePlanQueries")
+ public void testNoRuleSeesLargeSarg(String query) {
+ if (query.contains("partition_key")) {
+ // The guarded planner has no partition metadata for partition hints.
+ return;
+ }
+ QueryEnvironment guarded = getGuardedQueryEnvironment();
+ try (QueryEnvironment.CompiledQuery compiledQuery =
guarded.compile(query)) {
+ compiledQuery.planQuery(1);
+ }
+ }
+
+ /// The guard of [#testNoRuleSeesLargeSarg] fails when sealing is off.
+ @Test
+ public void testGuardFailsWithoutSealing() {
+ QueryEnvironment guarded = getGuardedQueryEnvironment();
+ Throwable error = expectThrows(Throwable.class, () -> {
+ try (QueryEnvironment.CompiledQuery compiledQuery =
+ guarded.compile("SET sealedInListThreshold=0; SELECT col1 FROM a
WHERE col3 IN (" + INTS + ")")) {
+ compiledQuery.planQuery(1);
+ }
+ });
+ while (error.getCause() != null && !(error instanceof AssertionError)) {
+ error = error.getCause();
+ }
+ assertTrue(error instanceof AssertionError &&
error.getMessage().contains("was offered SEARCH"),
+ String.valueOf(error));
+ }
+
+ private String explain(String query) {
+ return explain(_queryEnvironment, query);
+ }
+
+ private static String explain(QueryEnvironment queryEnvironment, String
query) {
+ return explain(queryEnvironment, query, "EXPLAIN PLAN FOR ");
+ }
+
+ private static String explain(QueryEnvironment queryEnvironment, String
query, String explainPrefix) {
+ int split = query.lastIndexOf(';') + 1;
+ String explainQuery = query.substring(0, split) + " " + explainPrefix +
query.substring(split).trim();
+ try (QueryEnvironment.CompiledQuery compiledQuery =
queryEnvironment.compile(explainQuery)) {
+ return compiledQuery.explain(1, null).getExplainPlan();
+ }
+ }
+
+ private String implementationPlan(String query) {
+ return explain(_queryEnvironment, query, "EXPLAIN IMPLEMENTATION PLAN FOR
").replaceAll("@localhost:\\d+",
+ "@host");
+ }
+
+ private List<String> serializedStages(String query) {
+ return serializedStages(_queryEnvironment, query);
+ }
+
+ private static List<String> serializedStages(QueryEnvironment
queryEnvironment, String query) {
+ try (QueryEnvironment.CompiledQuery compiledQuery =
queryEnvironment.compile(query)) {
+ DispatchableSubPlan plan = compiledQuery.planQuery(1).getQueryPlan();
+ List<String> stages = new ArrayList<>();
+ for (DispatchablePlanFragment fragment :
plan.getQueryStagesWithoutRoot()) {
+ // Also compare the segments of each worker.
+
stages.add(PlanNodeSerializer.process(fragment.getPlanFragment().getFragmentRoot())
+ "\n"
+ + fragment.getWorkerIdToSegmentsMap());
+ }
+ return stages;
+ }
+ }
+
+ private static QueryEnvironment getGuardedQueryEnvironment() {
+ return newQueryEnvironment(TABLE_SCHEMAS, true);
+ }
+
+ /// Returns a planner that is set up like the broker for queries without
null handling. With `guarded`, a
+ /// [LargeSargGuardRule] runs before the first rule and after each rule of
every phase.
+ private static QueryEnvironment newQueryEnvironment(Map<String, Schema>
schemas, boolean guarded) {
+ MockRoutingManagerFactory factory = new MockRoutingManagerFactory(1, 2);
+ for (Map.Entry<String, Schema> entry : schemas.entrySet()) {
+ factory.registerTable(entry.getValue(), entry.getKey());
+ }
+ for (Map.Entry<String, List<String>> entry : SERVER1_SEGMENTS.entrySet()) {
+ for (String segment : entry.getValue()) {
+ factory.registerSegment(1, entry.getKey(), segment);
+ }
+ }
+ for (Map.Entry<String, List<String>> entry : SERVER2_SEGMENTS.entrySet()) {
+ for (String segment : entry.getValue()) {
+ factory.registerSegment(2, entry.getKey(), segment);
+ }
+ }
+ RoutingManager routingManager = factory.buildRoutingManager(null);
+ TableCache tableCache = factory.buildTableCache();
+ RuleSetCustomizer guard = new RuleSetCustomizer() {
+ @Override
+ public void customize(Phase phase, List<RelOptRule> rules) {
+ List<RelOptRule> guarded = new ArrayList<>();
+ guarded.add(new LargeSargGuardRule(phase + "_guard_0"));
+ for (RelOptRule rule : rules) {
+ guarded.add(rule);
+ guarded.add(new LargeSargGuardRule(phase + "_guard_" +
guarded.size()));
+ }
+ rules.clear();
+ rules.addAll(guarded);
+ }
+ };
+ ImmutableQueryEnvironment.Config.Builder config =
QueryEnvironment.configBuilder()
+ .requestId(1L)
+ .database(CommonConstants.DEFAULT_DATABASE)
+ .tableCache(tableCache)
+ .workerManager(new WorkerManager("Broker_localhost", "localhost", 3,
routingManager))
+ .isNullHandlingEnabled(false)
+ .defaultSealedInListThreshold(GUARD_THRESHOLD);
+ if (guarded) {
+ config.ruleSet(new PinotRuleSet(List.of(new DefaultRuleSetCustomizer(),
guard)));
+ }
+ return new QueryEnvironment(config.build());
+ }
+
+ /// Never fires; fails when it is offered a node with an unsealed large
Sarg, or with an AND/OR of many comparisons of
+ /// one operand to literals.
+ private static final class LargeSargGuardRule extends RelOptRule {
+ LargeSargGuardRule(String description) {
+ // Deprecated operand API, like the other Pinot rules that match any
node: this rule never transforms.
+ super(operand(RelNode.class, any()), description);
+ }
+
+ @Override
+ public boolean matches(RelOptRuleCall call) {
+ RelNode rel = call.rel(0);
+ rel.accept(new RexShuttle() {
+ @Override
+ public RexNode visitCall(RexCall rexCall) {
+ int size = 0;
+ if (rexCall.getKind() == SqlKind.SEARCH) {
+ size = ((RexLiteral)
rexCall.getOperands().get(1)).getValueAs(Sarg.class).rangeSet.asRanges().size();
+ } else if (rexCall.getKind() == SqlKind.AND || rexCall.getKind() ==
SqlKind.OR) {
+ // Only comparisons of one operand can fold into one Sarg.
+ Map<RexNode, Integer> comparisons = new HashMap<>();
+ for (RexNode operand : rexCall.getOperands()) {
+ RexNode compared = comparedToLiteral(operand);
+ if (compared != null) {
+ size = Math.max(size, comparisons.merge(compared, 1,
Integer::sum));
+ }
+ }
+ }
+ if (size >= GUARD_THRESHOLD) {
+ throw new AssertionError(description + " was offered " +
rexCall.getKind() + " of size " + size + " in "
+ + rel);
+ }
+ return super.visitCall(rexCall);
+ }
+ });
+ return false;
+ }
+
+ @Override
+ public void onMatch(RelOptRuleCall call) {
+ }
+
+ @Nullable
+ private static RexNode comparedToLiteral(RexNode node) {
+ if (!node.isA(SqlKind.COMPARISON)) {
+ return null;
+ }
+ List<RexNode> operands = ((RexCall) node).getOperands();
+ if (operands.size() != 2) {
+ return null;
+ }
+ if (operands.get(1) instanceof RexLiteral) {
+ return operands.get(0);
+ }
+ return operands.get(0) instanceof RexLiteral ? operands.get(1) : null;
+ }
+ }
+}
diff --git
a/pinot-query-planner/src/test/java/org/apache/pinot/query/planner/logical/RexExpressionUtilsTest.java
b/pinot-query-planner/src/test/java/org/apache/pinot/query/planner/logical/RexExpressionUtilsTest.java
index a0883b349a2..70a7199f479 100644
---
a/pinot-query-planner/src/test/java/org/apache/pinot/query/planner/logical/RexExpressionUtilsTest.java
+++
b/pinot-query-planner/src/test/java/org/apache/pinot/query/planner/logical/RexExpressionUtilsTest.java
@@ -21,6 +21,7 @@ package org.apache.pinot.query.planner.logical;
import com.google.common.collect.ImmutableRangeSet;
import com.google.common.collect.Range;
import java.math.BigDecimal;
+import java.util.List;
import java.util.UUID;
import org.apache.calcite.rel.type.RelDataTypeFactory;
import org.apache.calcite.rex.RexBuilder;
@@ -585,4 +586,152 @@ public class RexExpressionUtilsTest {
Assert.assertTrue(result instanceof RexExpression.Literal);
Assert.assertEquals(result, RexExpression.Literal.TRUE);
}
+
+ /// `x IN (<20 values>) OR x > 100`: the points become one IN instead of one
pair of comparisons each.
+ @Test
+ public void testHandleSearchPointsAndRange() {
+ RexExpression result = convertIntSearch(RexUnknownAs.UNKNOWN,
+ withRanges(points(1, 20), Range.greaterThan(BigDecimal.valueOf(100))));
+ Assert.assertEquals(result,
+ call(SqlKind.OR, in(SqlKind.IN, 1, 20), call(SqlKind.GREATER_THAN,
ref(), intLiteral(100))));
+ }
+
+ /// `x NOT IN (<20 values>) AND x > 0`: the complement holds the points, so
they become one NOT_IN.
+ @Test
+ public void testHandleSearchComplementedPointsAndRange() {
+ ImmutableRangeSet<BigDecimal> rangeSet =
+ ImmutableRangeSet.copyOf(withRanges(points(1, 20),
Range.atMost(BigDecimal.ZERO)).complement());
+ RexExpression result = convertIntSearch(RexUnknownAs.UNKNOWN, rangeSet);
+ Assert.assertEquals(result, call(SqlKind.AND, in(SqlKind.NOT_IN, 1, 20),
+ call(SqlKind.GREATER_THAN, ref(), intLiteral(0))));
+ }
+
+ /// `x BETWEEN 0 AND 1000 AND x NOT IN (<25 values>)`: a complement range on
each side becomes one comparison each.
+ @Test
+ public void testHandleSearchComplementedPointsInBoundedRange() {
+ ImmutableRangeSet<BigDecimal> rangeSet =
ImmutableRangeSet.copyOf(withRanges(points(1, 25),
+ Range.lessThan(BigDecimal.ZERO),
Range.greaterThan(BigDecimal.valueOf(1000))).complement());
+ RexExpression result = convertIntSearch(RexUnknownAs.UNKNOWN, rangeSet);
+ Assert.assertEquals(result, call(SqlKind.AND, in(SqlKind.NOT_IN, 1, 25),
+ call(SqlKind.GREATER_THAN_OR_EQUAL, ref(), intLiteral(0)),
+ call(SqlKind.LESS_THAN_OR_EQUAL, ref(), intLiteral(1000))));
+ }
+
+ /// `(x IN (<20 values>) OR x BETWEEN 100 AND 200) OR x IS NULL`, and the
same with `AND x IS NOT NULL`.
+ @Test
+ public void testHandleSearchPointsAndClosedRangeWithNullCheck() {
+ ImmutableRangeSet<BigDecimal> rangeSet =
+ withRanges(points(1, 20), Range.closed(BigDecimal.valueOf(100),
BigDecimal.valueOf(200)));
+ RexExpression ranges = call(SqlKind.OR, in(SqlKind.IN, 1, 20),
+ call(SqlKind.AND, call(SqlKind.GREATER_THAN_OR_EQUAL, ref(),
intLiteral(100)),
+ call(SqlKind.LESS_THAN_OR_EQUAL, ref(), intLiteral(200))));
+ Assert.assertEquals(convertIntSearch(RexUnknownAs.TRUE, rangeSet),
+ call(SqlKind.OR, ranges, call(SqlKind.IS_NULL, ref())));
+ Assert.assertEquals(convertIntSearch(RexUnknownAs.FALSE, rangeSet),
+ call(SqlKind.AND, ranges, call(SqlKind.IS_NOT_NULL, ref())));
+ }
+
+ /// A few points next to a range keep one pair of comparisons per point, as
before.
+ @Test
+ public void testHandleSearchFewPointsAndRange() {
+ RexExpression result = convertIntSearch(RexUnknownAs.UNKNOWN,
+ withRanges(points(1, RexExpressionUtils.MIN_POINTS_FOR_IN - 1),
Range.greaterThan(BigDecimal.valueOf(100))));
+ Assert.assertTrue(result instanceof RexExpression.FunctionCall);
+ RexExpression.FunctionCall or = (RexExpression.FunctionCall) result;
+ Assert.assertEquals(or.getFunctionName(), SqlKind.OR.name());
+ Assert.assertEquals(or.getFunctionOperands().size(),
RexExpressionUtils.MIN_POINTS_FOR_IN);
+ Assert.assertEquals(or.getFunctionOperands().get(0), call(SqlKind.AND,
+ call(SqlKind.GREATER_THAN_OR_EQUAL, ref(), intLiteral(1)),
+ call(SqlKind.LESS_THAN_OR_EQUAL, ref(), intLiteral(1))));
+ }
+
+ /// `x BETWEEN 0 AND 10 AND x NOT IN (3, 5)`: no points and 2 complement
points. This must not become an IN without
+ /// values.
+ @Test
+ public void testHandleSearchRangesWithoutPoints() {
+ RexExpression result = convertIntSearch(RexUnknownAs.UNKNOWN, withRanges(
+ Range.closedOpen(BigDecimal.ZERO, BigDecimal.valueOf(3)),
+ Range.open(BigDecimal.valueOf(3), BigDecimal.valueOf(5)),
+ Range.openClosed(BigDecimal.valueOf(5), BigDecimal.valueOf(10))));
+ Assert.assertEquals(result, call(SqlKind.OR,
+ call(SqlKind.AND, call(SqlKind.GREATER_THAN_OR_EQUAL, ref(),
intLiteral(0)),
+ call(SqlKind.LESS_THAN, ref(), intLiteral(3))),
+ call(SqlKind.AND, call(SqlKind.GREATER_THAN, ref(), intLiteral(3)),
+ call(SqlKind.LESS_THAN, ref(), intLiteral(5))),
+ call(SqlKind.AND, call(SqlKind.GREATER_THAN, ref(), intLiteral(5)),
+ call(SqlKind.LESS_THAN_OR_EQUAL, ref(), intLiteral(10)))));
+ }
+
+ /// BIG_DECIMAL points next to a range stay comparisons, because
intermediate stages match IN values with `equals`.
+ @Test
+ public void testHandleSearchBigDecimalPointsAndRange() {
+ RexInputRef inputRef =
_rexBuilder.makeInputRef(_typeFactory.createSqlType(SqlTypeName.DECIMAL, 38,
2), 0);
+ Sarg<BigDecimal> sarg = Sarg.of(RexUnknownAs.UNKNOWN,
+ withRanges(points(1, 25), Range.greaterThan(BigDecimal.valueOf(100))));
+ RexLiteral searchLiteral =
+ _rexBuilder.makeSearchArgumentLiteral(sarg,
_typeFactory.createSqlType(SqlTypeName.DECIMAL, 38, 2));
+ RexExpression result = RexExpressionUtils.fromRexCall(
+ (RexCall) _rexBuilder.makeCall(SqlStdOperatorTable.SEARCH, inputRef,
searchLiteral));
+ Assert.assertTrue(result instanceof RexExpression.FunctionCall);
+ RexExpression.FunctionCall or = (RexExpression.FunctionCall) result;
+ Assert.assertEquals(or.getFunctionName(), SqlKind.OR.name());
+ Assert.assertEquals(or.getFunctionOperands().size(), 26);
+ for (RexExpression operand : or.getFunctionOperands()) {
+ Assert.assertNotEquals(((RexExpression.FunctionCall)
operand).getFunctionName(), SqlKind.IN.name());
+ }
+ }
+
+ /// Point ranges for the values `from` to `to`, both included.
+ private static ImmutableRangeSet<BigDecimal> points(int from, int to) {
+ ImmutableRangeSet.Builder<BigDecimal> builder =
ImmutableRangeSet.builder();
+ for (int i = from; i <= to; i++) {
+ builder.add(Range.singleton(BigDecimal.valueOf(i)));
+ }
+ return builder.build();
+ }
+
+ @SafeVarargs
+ private static ImmutableRangeSet<BigDecimal> withRanges(Range<BigDecimal>...
ranges) {
+ return withRanges(ImmutableRangeSet.of(), ranges);
+ }
+
+ @SafeVarargs
+ private static ImmutableRangeSet<BigDecimal>
withRanges(ImmutableRangeSet<BigDecimal> rangeSet,
+ Range<BigDecimal>... ranges) {
+ ImmutableRangeSet.Builder<BigDecimal> builder =
ImmutableRangeSet.<BigDecimal>builder().addAll(rangeSet);
+ for (Range<BigDecimal> range : ranges) {
+ builder.add(range);
+ }
+ return builder.build();
+ }
+
+ private static RexExpression in(SqlKind kind, int from, int to) {
+ RexExpression[] operands = new RexExpression[to - from + 2];
+ operands[0] = ref();
+ for (int i = from; i <= to; i++) {
+ operands[i - from + 1] = intLiteral(i);
+ }
+ return call(kind, operands);
+ }
+
+ private RexExpression convertIntSearch(RexUnknownAs nullAs,
ImmutableRangeSet<BigDecimal> rangeSet) {
+ RexInputRef inputRef =
_rexBuilder.makeInputRef(_typeFactory.createSqlType(SqlTypeName.INTEGER), 0);
+ Sarg<BigDecimal> sarg = Sarg.of(nullAs, rangeSet);
+ RexLiteral searchLiteral =
+ _rexBuilder.makeSearchArgumentLiteral(sarg,
_typeFactory.createSqlType(SqlTypeName.INTEGER));
+ return RexExpressionUtils.fromRexCall(
+ (RexCall) _rexBuilder.makeCall(SqlStdOperatorTable.SEARCH, inputRef,
searchLiteral));
+ }
+
+ private static RexExpression ref() {
+ return new RexExpression.InputRef(0);
+ }
+
+ private static RexExpression intLiteral(int value) {
+ return new RexExpression.Literal(ColumnDataType.INT, value);
+ }
+
+ private static RexExpression call(SqlKind kind, RexExpression... operands) {
+ return new RexExpression.FunctionCall(ColumnDataType.BOOLEAN, kind.name(),
List.of(operands));
+ }
}
diff --git a/pinot-query-runtime/src/test/resources/queries/LargeInLists.json
b/pinot-query-runtime/src/test/resources/queries/LargeInLists.json
new file mode 100644
index 00000000000..0c73888acb4
--- /dev/null
+++ b/pinot-query-runtime/src/test/resources/queries/LargeInLists.json
@@ -0,0 +1,602 @@
+{
+ "large_in_lists": {
+ "comment": "IN lists with at least 20 values, which the multi-stage
planner seals during optimization",
+ "tables": {
+ "facts": {
+ "schema": [
+ {
+ "name": "id",
+ "type": "INT"
+ },
+ {
+ "name": "grp",
+ "type": "STRING"
+ },
+ {
+ "name": "val",
+ "type": "LONG"
+ },
+ {
+ "name": "amt",
+ "type": "DOUBLE"
+ },
+ {
+ "name": "dim_id",
+ "type": "INT"
+ }
+ ],
+ "inputs": [
+ [
+ 1,
+ "g1",
+ 3,
+ 1.5,
+ 2
+ ],
+ [
+ 2,
+ "g2",
+ 6,
+ 2.5,
+ 3
+ ],
+ [
+ 3,
+ "g3",
+ 9,
+ 3.5,
+ 4
+ ],
+ [
+ 4,
+ "g4",
+ 12,
+ 4.5,
+ 5
+ ],
+ [
+ 5,
+ "g5",
+ 15,
+ 5.5,
+ 6
+ ],
+ [
+ 6,
+ "g6",
+ 18,
+ 6.5,
+ 7
+ ],
+ [
+ 7,
+ "g7",
+ 21,
+ 7.5,
+ 8
+ ],
+ [
+ 8,
+ "g8",
+ 24,
+ 8.5,
+ 9
+ ],
+ [
+ 9,
+ "g9",
+ 27,
+ 9.5,
+ 10
+ ],
+ [
+ 10,
+ "g10",
+ 30,
+ 10.5,
+ 11
+ ],
+ [
+ 11,
+ "g11",
+ 33,
+ 11.5,
+ 12
+ ],
+ [
+ 12,
+ "g12",
+ 36,
+ 12.5,
+ 13
+ ],
+ [
+ 13,
+ "g13",
+ 39,
+ 13.5,
+ 14
+ ],
+ [
+ 14,
+ "g14",
+ 42,
+ 14.5,
+ 15
+ ],
+ [
+ 15,
+ "g15",
+ 45,
+ 15.5,
+ 16
+ ],
+ [
+ 16,
+ "g16",
+ 48,
+ 16.5,
+ 17
+ ],
+ [
+ 17,
+ "g17",
+ 51,
+ 17.5,
+ 1
+ ],
+ [
+ 18,
+ "g18",
+ 54,
+ 18.5,
+ 2
+ ],
+ [
+ 19,
+ "g19",
+ 57,
+ 19.5,
+ 3
+ ],
+ [
+ 20,
+ "g20",
+ 60,
+ 20.5,
+ 4
+ ],
+ [
+ 21,
+ "g21",
+ 63,
+ 21.5,
+ 5
+ ],
+ [
+ 22,
+ "g22",
+ 66,
+ 22.5,
+ 6
+ ],
+ [
+ 23,
+ "g23",
+ 69,
+ 23.5,
+ 7
+ ],
+ [
+ 24,
+ "g24",
+ 72,
+ 24.5,
+ 8
+ ],
+ [
+ 25,
+ "g25",
+ 75,
+ 25.5,
+ 9
+ ],
+ [
+ 26,
+ "g26",
+ 78,
+ 26.5,
+ 10
+ ],
+ [
+ 27,
+ "g27",
+ 81,
+ 27.5,
+ 11
+ ],
+ [
+ 28,
+ "g28",
+ 84,
+ 28.5,
+ 12
+ ],
+ [
+ 29,
+ "g29",
+ 87,
+ 29.5,
+ 13
+ ],
+ [
+ 30,
+ "g0",
+ 90,
+ 30.5,
+ 14
+ ],
+ [
+ 31,
+ "g1",
+ 93,
+ 31.5,
+ 15
+ ],
+ [
+ 32,
+ "g2",
+ 96,
+ 32.5,
+ 16
+ ],
+ [
+ 33,
+ "g3",
+ 99,
+ 33.5,
+ 17
+ ],
+ [
+ 34,
+ "g4",
+ 102,
+ 34.5,
+ 1
+ ],
+ [
+ 35,
+ "g5",
+ 105,
+ 35.5,
+ 2
+ ],
+ [
+ 36,
+ "g6",
+ 108,
+ 36.5,
+ 3
+ ],
+ [
+ 37,
+ "g7",
+ 111,
+ 37.5,
+ 4
+ ],
+ [
+ 38,
+ "g8",
+ 114,
+ 38.5,
+ 5
+ ],
+ [
+ 39,
+ "g9",
+ 117,
+ 39.5,
+ 6
+ ],
+ [
+ 40,
+ "g10",
+ 120,
+ 40.5,
+ 7
+ ]
+ ]
+ },
+ "dims": {
+ "schema": [
+ {
+ "name": "dim_id",
+ "type": "INT"
+ },
+ {
+ "name": "name",
+ "type": "STRING"
+ },
+ {
+ "name": "cat",
+ "type": "INT"
+ }
+ ],
+ "inputs": [
+ [
+ 1,
+ "name1",
+ 2
+ ],
+ [
+ 2,
+ "name2",
+ 4
+ ],
+ [
+ 3,
+ "name3",
+ 6
+ ],
+ [
+ 4,
+ "name4",
+ 8
+ ],
+ [
+ 5,
+ "name5",
+ 10
+ ],
+ [
+ 6,
+ "name6",
+ 12
+ ],
+ [
+ 7,
+ "name7",
+ 14
+ ],
+ [
+ 8,
+ "name8",
+ 16
+ ],
+ [
+ 9,
+ "name9",
+ 18
+ ],
+ [
+ 10,
+ "name10",
+ 20
+ ],
+ [
+ 11,
+ "name11",
+ 22
+ ],
+ [
+ 12,
+ "name12",
+ 24
+ ],
+ [
+ 13,
+ "name13",
+ 26
+ ],
+ [
+ 14,
+ "name14",
+ 28
+ ],
+ [
+ 15,
+ "name15",
+ 30
+ ],
+ [
+ 16,
+ "name16",
+ 32
+ ],
+ [
+ 17,
+ "name17",
+ 34
+ ],
+ [
+ 18,
+ "name18",
+ 36
+ ]
+ ]
+ }
+ },
+ "queries": [
+ {
+ "description": "where in",
+ "sql": "SELECT id, grp FROM {facts} WHERE id IN (1, 3, 5, 7, 9, 11,
13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45, 47, 49, 51,
53, 55, 57, 59)"
+ },
+ {
+ "description": "where not in",
+ "sql": "SELECT COUNT(*) FROM {facts} WHERE id NOT IN (1, 3, 5, 7, 9,
11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45, 47, 49,
51, 53, 55, 57, 59)"
+ },
+ {
+ "description": "in or range on another column",
+ "sql": "SELECT COUNT(*) FROM {facts} WHERE id IN (1, 3, 5, 7, 9, 11,
13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45, 47, 49, 51,
53, 55, 57, 59) OR val > 90"
+ },
+ {
+ "description": "in or range on the same column",
+ "sql": "SELECT id FROM {facts} WHERE id IN (1, 3, 5, 7, 9, 11, 13, 15,
17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45, 47, 49, 51, 53, 55,
57, 59) OR id > 35"
+ },
+ {
+ "description": "not in and range on the same column",
+ "sql": "SELECT id FROM {facts} WHERE id NOT IN (1, 3, 5, 7, 9, 11, 13,
15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45, 47, 49, 51, 53,
55, 57, 59) AND id > 5"
+ },
+ {
+ "description": "in or between on the same column",
+ "sql": "SELECT id FROM {facts} WHERE id IN (1, 3, 5, 7, 9, 11, 13, 15,
17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45, 47, 49, 51, 53, 55,
57, 59) OR id BETWEEN 10 AND 14"
+ },
+ {
+ "description": "case and filter",
+ "sql": "SELECT SUM(CASE WHEN id IN (1, 3, 5, 7, 9, 11, 13, 15, 17, 19,
21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45, 47, 49, 51, 53, 55, 57, 59)
THEN 1 ELSE 0 END), COUNT(*) FILTER (WHERE id NOT IN (2, 5, 8, 11, 14, 17, 20,
23, 26, 29, 32, 35, 38, 41, 44, 47, 50, 53, 56, 59)) FROM {facts}"
+ },
+ {
+ "description": "select list",
+ "sql": "SELECT id, id IN (1, 3, 5, 7, 9, 11, 13, 15, 17, 19, 21, 23,
25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45, 47, 49, 51, 53, 55, 57, 59) FROM
{facts}"
+ },
+ {
+ "description": "group by case",
+ "sql": "SELECT CASE WHEN id IN (1, 3, 5, 7, 9, 11, 13, 15, 17, 19, 21,
23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45, 47, 49, 51, 53, 55, 57, 59)
THEN 'odd' ELSE 'even' END AS k, COUNT(*) FROM {facts} GROUP BY CASE WHEN id IN
(1, 3, 5, 7, 9, 11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41,
43, 45, 47, 49, 51, 53, 55, 57, 59) THEN 'odd' ELSE 'even' END"
+ },
+ {
+ "description": "having",
+ "sql": "SELECT grp, COUNT(*) FROM {facts} GROUP BY grp HAVING SUM(val)
IN (2, 5, 8, 11, 14, 17, 20, 23, 26, 29, 32, 35, 38, 41, 44, 47, 50, 53, 56,
59, 3, 6, 9, 36, 99)"
+ },
+ {
+ "description": "join on",
+ "sql": "SELECT f.id, d.name FROM {facts} f JOIN {dims} d ON f.dim_id =
d.dim_id AND d.cat IN (2, 5, 8, 11, 14, 17, 20, 23, 26, 29, 32, 35, 38, 41, 44,
47, 50, 53, 56, 59)"
+ },
+ {
+ "description": "left join on left column",
+ "sql": "SELECT f.id, d.name FROM {facts} f LEFT JOIN {dims} d ON
f.dim_id = d.dim_id AND f.id IN (1, 3, 5, 7, 9, 11, 13, 15, 17, 19, 21, 23, 25,
27, 29, 31, 33, 35, 37, 39, 41, 43, 45, 47, 49, 51, 53, 55, 57, 59)"
+ },
+ {
+ "description": "left join where right in",
+ "sql": "SELECT f.id, d.name FROM {facts} f LEFT JOIN {dims} d ON
f.dim_id = d.dim_id WHERE d.cat IN (2, 5, 8, 11, 14, 17, 20, 23, 26, 29, 32,
35, 38, 41, 44, 47, 50, 53, 56, 59)"
+ },
+ {
+ "description": "left join where right not in or null",
+ "sql": "SELECT f.id, d.name FROM {facts} f LEFT JOIN {dims} d ON
f.dim_id = d.dim_id AND d.cat < 20 WHERE d.cat NOT IN (2, 5, 8, 11, 14, 17, 20,
23, 26, 29, 32, 35, 38, 41, 44, 47, 50, 53, 56, 59) OR d.cat IS NULL"
+ },
+ {
+ "description": "join key",
+ "sql": "SELECT f.id, d.name FROM {facts} f JOIN {dims} d ON f.dim_id =
d.dim_id WHERE f.dim_id IN (1, 3, 5, 7, 9, 11, 13, 15, 17, 19, 21, 23, 25, 27,
29, 31, 33, 35, 37, 39, 41, 43, 45, 47, 49, 51, 53, 55, 57, 59)"
+ },
+ {
+ "description": "strings",
+ "sql": "SELECT id FROM {facts} WHERE grp IN ('g0', 'g2', 'g4', 'g6',
'g8', 'g10', 'g12', 'g14', 'g16', 'g18', 'g20', 'g22', 'g24', 'g26', 'g28',
'g30', 'g32', 'g34', 'g36', 'g38', 'g40', 'g42', 'g44', 'g46', 'g48')"
+ },
+ {
+ "description": "expression",
+ "sql": "SELECT id FROM {facts} WHERE UPPER(grp) IN ('G0', 'G2', 'G4',
'G6', 'G8', 'G10', 'G12', 'G14', 'G16', 'G18', 'G20', 'G22', 'G24', 'G26',
'G28', 'G30', 'G32', 'G34', 'G36', 'G38', 'G40', 'G42', 'G44', 'G46', 'G48')"
+ },
+ {
+ "description": "cte twice",
+ "sql": "WITH t AS (SELECT id, dim_id FROM {facts} WHERE id IN (1, 3,
5, 7, 9, 11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43,
45, 47, 49, 51, 53, 55, 57, 59)) SELECT COUNT(*) FROM t t1 JOIN t t2 ON
t1.dim_id = t2.dim_id"
+ },
+ {
+ "description": "union all",
+ "sql": "SELECT id FROM {facts} WHERE id IN (1, 3, 5, 7, 9, 11, 13, 15,
17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45, 47, 49, 51, 53, 55,
57, 59) UNION ALL SELECT dim_id FROM {dims} WHERE cat NOT IN (2, 5, 8, 11, 14,
17, 20, 23, 26, 29, 32, 35, 38, 41, 44, 47, 50, 53, 56, 59)"
+ },
+ {
+ "description": "window",
+ "sql": "SELECT id, ROW_NUMBER() OVER (PARTITION BY CASE WHEN id IN (1,
3, 5, 7, 9, 11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43,
45, 47, 49, 51, 53, 55, 57, 59) THEN 1 ELSE 0 END ORDER BY id) FROM {facts}"
+ },
+ {
+ "description": "null in list",
+ "sql": "SELECT COUNT(*) FROM {facts} WHERE id IN (1, 3, 5, 7, 9, 11,
13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45, 47, 49, 51,
53, 55, 57, 59, NULL)"
+ },
+ {
+ "description": "not in with null in list",
+ "sql": "SELECT COUNT(*) FROM {facts} WHERE id NOT IN (1, 3, 5, 7, 9,
11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45, 47, 49,
51, 53, 55, 57, 59, NULL)"
+ },
+ {
+ "description": "column in list",
+ "sql": "SELECT COUNT(*) FROM {facts} WHERE id IN (dim_id, 1, 3, 5, 7,
9, 11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45, 47,
49, 51, 53, 55, 57, 59)"
+ },
+ {
+ "description": "long column",
+ "sql": "SELECT COUNT(*) FROM {facts} WHERE val IN (1, 3, 5, 7, 9, 11,
13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45, 47, 49, 51,
53, 55, 57, 59)"
+ },
+ {
+ "description": "double column",
+ "sql": "SELECT COUNT(*) FROM {facts} WHERE amt IN (1.5, 3.5, 5.5, 7.5,
9.5, 11.5, 13.5, 15.5, 17.5, 19.5, 21.5, 23.5, 25.5, 27.5, 29.5, 31.5, 33.5,
35.5, 37.5, 39.5, 41.5, 43.5, 45.5, 47.5, 49.5)"
+ },
+ {
+ "description": "equality and list",
+ "sql": "SELECT COUNT(*) FROM {facts} WHERE id = 3 AND id IN (1, 3, 5,
7, 9, 11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45,
47, 49, 51, 53, 55, 57, 59)"
+ },
+ {
+ "description": "contradicting equality and list",
+ "sql": "SELECT COUNT(*) FROM {facts} WHERE id = 4 AND id IN (1, 3, 5,
7, 9, 11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45,
47, 49, 51, 53, 55, 57, 59)"
+ },
+ {
+ "description": "literal operand",
+ "sql": "SELECT COUNT(*) FROM (SELECT 5 AS x, id FROM {facts}) t WHERE
t.x IN (1, 3, 5, 7, 9, 11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37,
39, 41, 43, 45, 47, 49, 51, 53, 55, 57, 59)"
+ },
+ {
+ "description": "in sub-query",
+ "sql": "SELECT id FROM {facts} WHERE dim_id IN (SELECT dim_id FROM
{dims} WHERE cat IN (2, 5, 8, 11, 14, 17, 20, 23, 26, 29, 32, 35, 38, 41, 44,
47, 50, 53, 56, 59))"
+ },
+ {
+ "description": "correlated exists",
+ "sql": "SELECT COUNT(*) FROM {facts} f WHERE EXISTS (SELECT 1 FROM
{dims} d WHERE d.dim_id = f.dim_id AND d.cat IN (2, 5, 8, 11, 14, 17, 20, 23,
26, 29, 32, 35, 38, 41, 44, 47, 50, 53, 56, 59))"
+ },
+ {
+ "description": "two lists on one column",
+ "sql": "SELECT id FROM {facts} WHERE id IN (1, 3, 5, 7, 9, 11, 13, 15,
17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45, 47, 49, 51, 53, 55,
57, 59) AND id IN (2, 5, 8, 11, 14, 17, 20, 23, 26, 29, 32, 35, 38, 41, 44, 47,
50, 53, 56, 59)"
+ },
+ {
+ "description": "user written or chain",
+ "sql": "SELECT id FROM {facts} WHERE id = 1 OR id = 3 OR id = 5 OR id
= 7 OR id = 9 OR id = 11 OR id = 13 OR id = 15 OR id = 17 OR id = 19 OR id = 21
OR id = 23 OR id = 25 OR id = 27 OR id = 29 OR id = 31 OR id = 33 OR id = 35 OR
id = 37 OR id = 39 OR id = 41 OR id = 43 OR id = 45 OR id = 47 OR id = 49 OR id
= 51 OR id = 53 OR id = 55 OR id = 57 OR id = 59"
+ },
+ {
+ "description": "null handling: not in over a nullable column",
+ "sql": "SET enableNullHandling = TRUE; WITH t AS (SELECT id, CASE WHEN
MOD(id, 5) = 0 THEN NULL ELSE id END AS nid FROM {facts}) SELECT COUNT(*) FROM
t WHERE nid NOT IN (1, 3, 5, 7, 9, 11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31,
33, 35, 37, 39, 41, 43, 45, 47, 49, 51, 53, 55, 57, 59)",
+ "h2Sql": "WITH t AS (SELECT id, CASE WHEN MOD(id, 5) = 0 THEN NULL
ELSE id END AS nid FROM {facts}) SELECT COUNT(*) FROM t WHERE nid NOT IN (1, 3,
5, 7, 9, 11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43,
45, 47, 49, 51, 53, 55, 57, 59)"
+ },
+ {
+ "description": "null handling: in over a nullable column in the select
list",
+ "sql": "SET enableNullHandling = TRUE; WITH t AS (SELECT id, CASE WHEN
MOD(id, 5) = 0 THEN NULL ELSE id END AS nid FROM {facts}) SELECT id, nid IN (1,
3, 5, 7, 9, 11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43,
45, 47, 49, 51, 53, 55, 57, 59) FROM t",
+ "h2Sql": "WITH t AS (SELECT id, CASE WHEN MOD(id, 5) = 0 THEN NULL
ELSE id END AS nid FROM {facts}) SELECT id, nid IN (1, 3, 5, 7, 9, 11, 13, 15,
17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45, 47, 49, 51, 53, 55,
57, 59) FROM t"
+ },
+ {
+ "description": "null handling: not in with null in the list, in the
select list",
+ "sql": "SET enableNullHandling = TRUE; WITH t AS (SELECT id, CASE WHEN
MOD(id, 5) = 0 THEN NULL ELSE id END AS nid FROM {facts}) SELECT id, nid NOT IN
(1, 3, 5, 7, 9, 11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41,
43, 45, 47, 49, 51, 53, 55, 57, 59, NULL) FROM t",
+ "h2Sql": "WITH t AS (SELECT id, CASE WHEN MOD(id, 5) = 0 THEN NULL
ELSE id END AS nid FROM {facts}) SELECT id, nid NOT IN (1, 3, 5, 7, 9, 11, 13,
15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45, 47, 49, 51, 53,
55, 57, 59, NULL) FROM t"
+ },
+ {
+ "description": "null handling: in result is null",
+ "sql": "SET enableNullHandling = TRUE; WITH t AS (SELECT id, CASE WHEN
MOD(id, 5) = 0 THEN NULL ELSE id END AS nid FROM {facts}) SELECT COUNT(*) FROM
t WHERE (nid IN (1, 3, 5, 7, 9, 11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33,
35, 37, 39, 41, 43, 45, 47, 49, 51, 53, 55, 57, 59)) IS NULL",
+ "h2Sql": "WITH t AS (SELECT id, CASE WHEN MOD(id, 5) = 0 THEN NULL
ELSE id END AS nid FROM {facts}) SELECT COUNT(*) FROM t WHERE (nid IN (1, 3, 5,
7, 9, 11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45,
47, 49, 51, 53, 55, 57, 59)) IS NULL"
+ },
+ {
+ "description": "null handling: in or is null over a nullable column",
+ "sql": "SET enableNullHandling = TRUE; WITH t AS (SELECT id, CASE WHEN
MOD(id, 5) = 0 THEN NULL ELSE id END AS nid FROM {facts}) SELECT COUNT(*) FROM
t WHERE nid IN (1, 3, 5, 7, 9, 11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33,
35, 37, 39, 41, 43, 45, 47, 49, 51, 53, 55, 57, 59) OR nid IS NULL",
+ "h2Sql": "WITH t AS (SELECT id, CASE WHEN MOD(id, 5) = 0 THEN NULL
ELSE id END AS nid FROM {facts}) SELECT COUNT(*) FROM t WHERE nid IN (1, 3, 5,
7, 9, 11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45,
47, 49, 51, 53, 55, 57, 59) OR nid IS NULL"
+ },
+ {
+ "description": "null handling: left join where right column in list",
+ "sql": "SET enableNullHandling = TRUE; SELECT f.id, d.name FROM
{facts} f LEFT JOIN {dims} d ON f.dim_id = d.dim_id AND d.cat < 20 WHERE d.cat
IN (2, 4, 6, 8, 10, 12, 14, 16, 18, 20, 22, 24, 26, 28, 30, 32, 34, 36, 38,
40)",
+ "h2Sql": "SELECT f.id, d.name FROM {facts} f LEFT JOIN {dims} d ON
f.dim_id = d.dim_id AND d.cat < 20 WHERE d.cat IN (2, 4, 6, 8, 10, 12, 14, 16,
18, 20, 22, 24, 26, 28, 30, 32, 34, 36, 38, 40)"
+ },
+ {
+ "description": "null handling: left join where right column not in
list",
+ "sql": "SET enableNullHandling = TRUE; SELECT f.id, d.name FROM
{facts} f LEFT JOIN {dims} d ON f.dim_id = d.dim_id AND d.cat < 20 WHERE d.cat
NOT IN (1, 3, 5, 7, 9, 11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37,
39)",
+ "h2Sql": "SELECT f.id, d.name FROM {facts} f LEFT JOIN {dims} d ON
f.dim_id = d.dim_id AND d.cat < 20 WHERE d.cat NOT IN (1, 3, 5, 7, 9, 11, 13,
15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39)"
+ },
+ {
+ "description": "mixed sarg without points in the select list",
+ "sql": "SELECT id, CASE WHEN id BETWEEN 0 AND 10 AND id NOT IN (3, 5)
THEN 1 ELSE 0 END FROM {facts}"
+ },
+ {
+ "description": "mixed sarg without points in having",
+ "sql": "SELECT grp, COUNT(*) FROM {facts} GROUP BY grp HAVING COUNT(*)
BETWEEN 0 AND 10 AND COUNT(*) NOT IN (3, 5)"
+ },
+ {
+ "description": "large list or range above a limit",
+ "sql": "SELECT id FROM (SELECT id, val FROM {facts} ORDER BY id LIMIT
30) t WHERE id IN (1, 3, 5, 7, 9, 11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31,
33, 35, 37, 39, 41, 43, 45, 47, 49, 51, 53, 55, 57, 59) OR id > 100"
+ },
+ {
+ "description": "large list and range on an aggregate",
+ "sql": "SELECT grp, SUM(val) FROM {facts} GROUP BY grp HAVING SUM(val)
NOT IN (1, 3, 5, 7, 9, 11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37,
39, 41, 43, 45, 47, 49, 51, 53, 55, 57, 59) AND SUM(val) > 0"
+ }
+ ]
+ }
+}
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 1d1d36f7986..7ab92db3953 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
@@ -798,6 +798,26 @@ public class CommonConstants {
// TODO: Change this default to something very high, as this
_optimnization_ is usually not beneficial.
public static final int DEFAULT_SORT_EXCHANGE_COPY_THRESHOLD = 10_000;
+ /// Config for the smallest IN list that the multi-stage planner hides
from Calcite's optimizer.
+ ///
+ /// Calcite keeps an IN list as one `SEARCH` call over a sorted range set.
Many planner rules and metadata
+ /// handlers rebuild that range set every time they touch the predicate,
so planning time grows with the list size
+ /// times the number of touches. IN lists (and `OR` chains of equalities
in a filter) with at least this many values
+ /// are sealed into an opaque predicate during optimization and restored
afterwards.
+ ///
+ /// Predicates in the same filter or join condition still fold into the
list first, as without sealing. Sealed lists
+ /// lose Calcite's value-level reasoning across plan nodes: a predicate
that a rule moves next to a sealed list on
+ /// the same column is not merged into it. Null checks keep their meaning:
without null handling, Calcite cannot see
+ /// into Pinot's `IS NULL` and `IS NOT NULL` operators, and with null
handling the servers apply SQL null semantics.
+ ///
+ /// A value of 0 or less disables sealing. Planning outside the broker's
multi-stage request handler (for example
+ /// for the controller `/sql` endpoint) does not read this broker config
and uses the default. The query option
+ /// works everywhere.
+ public static final String CONFIG_OF_SEALED_IN_LIST_THRESHOLD =
"pinot.broker.multistage.sealed.in.list.threshold";
+ /// Same as the default of Calcite's
`SqlToRelConverter.Config#getInSubQueryThreshold()`: the size from which
+ /// Calcite itself stops treating an IN list as a scalar predicate.
+ public static final int DEFAULT_SEALED_IN_LIST_THRESHOLD = 20;
+
public static class Request {
public static final String SQL = "sql";
public static final String SQL_V1 = "sqlV1";
@@ -1148,6 +1168,9 @@ public class CommonConstants {
/// Option to customize the value of
[Broker#CONFIG_OF_SORT_EXCHANGE_COPY_THRESHOLD]
public static final String SORT_EXCHANGE_COPY_THRESHOLD =
"sortExchangeCopyThreshold";
+ /// Option to customize the value of
[Broker#CONFIG_OF_SEALED_IN_LIST_THRESHOLD]
+ public static final String SEALED_IN_LIST_THRESHOLD =
"sealedInListThreshold";
+
// Vector search query options
/// Number of inverted-list probes for IVF-based vector indexes.
Higher values improve recall
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]