This is an automated email from the ASF dual-hosted git repository.
englefly pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new 1ff44a04bc8 [fix](rf) Push runtime filters into a shared CTE producer
per filter identity (#68340)
1ff44a04bc8 is described below
commit 1ff44a04bc8c5b5bc385bf5d367f4676ad4a4191
Author: minghong <[email protected]>
AuthorDate: Fri Oct 9 16:02:39 2026 +0800
[fix](rf) Push runtime filters into a shared CTE producer per filter
identity (#68340)
### What problem does this PR solve?
Problem Summary:
A query that references a materialized CTE more than once and whose
references need different row
ranges silently loses rows when runtime filters are enabled: the runtime
filters generated for one
reference are pushed into the shared CTE producer scan and prune the
rows that the other references
still need. The result is empty or shorter than expected.
Reproduction on internal OLAP tables (no external table needed):
CREATE TABLE f (k INT) ... ; INSERT INTO f VALUES (0), (1), (5), (10);
CREATE TABLE b (x INT) ... ; INSERT INTO b VALUES (3);
SET enable_cte_materialize=true;
SET inline_cte_referenced_threshold=0;
SET runtime_filter_type='MIN_MAX';
SET runtime_filter_mode='GLOBAL';
WITH t AS (SELECT k, ABS(k) AS v FROM f)
SELECT c2.k AS lo, c1.k AS hi, c2.v AS lo_v, c1.v AS hi_v, b.x
FROM t c2 CROSS JOIN t c1 CROSS JOIN b
WHERE c1.k > b.x AND c2.k < b.x
ORDER BY lo, hi;
The query returns no row at all, while `runtime_filter_mode='OFF'`
returns the four rows
`(0,5) (0,10) (1,5) (1,10)`. The plan shows both consumer-side filters
installed on the shared
producer scan:
0:VOlapScanNode ... runtime filters: RF003[min] -> k, RF004[max] -> k
Root cause in code:
`RuntimeFilterGenerator` moved the runtime filters of the CTE consumers
into their shared producer
whenever the consumers' `srcExpr` sets intersect, and removed them from
the consumers. The code that
moved them only checked that the filters map to the same producer target
expression, not that they are
the same filter. Here `c1.k > b.x` generates a
`MIN_MAX/MIN` filter (k >= min(x)) while `c2.k < b.x` generates a
`MIN_MAX/MAX` filter (k <=
max(x)) on the same producer column, so both were installed on the
shared scan and their
conjunction (`k >= 3 AND k <= 3`) pruned every row.
Fix:
Decide per filter identity instead of per source expression:
`pushRuntimeFiltersIntoCTEProducer()`
collects the runtime filters of the CTE consumers, and
`selectPushableRuntimeFilters()` selects the
groups that may be pushed. The filters of one source expression
are grouped by identity -- same `TRuntimeFilterType`, same
`TMinMaxRuntimeFilterType` and same
target expression on the producer, that is the same predicate on the
producer's rows -- and each
group is pushed on its own when every consumer of the CTE applies a
filter of that identity. Only
filters that all consumers would apply anyway are applied on the shared
producer, which keeps the
pushdown equivalent to the per-consumer filters:
- `t c1 where c1.k > b.x` and `t c2 where c2.k < b.x`: the MIN group is
applied by one consumer and
the MAX group by the other, so neither group is pushed and the bug is
fixed;
- every consumer applying the same MIN_MAX and IN_OR_BLOOM filters,
which is the common shape when
`runtime_filter_type` enables both (the default value 12): both groups
are still pushed, so the
optimization is kept and TPC-DS q95 keeps the runtime filters on the
inner scan of its CTE body;
- a filter that only some consumers apply is skipped while the groups
the other consumers apply
are still pushed;
- filters which differ in the NULL semantics are never substituted for
each other, so a null aware
consumer keeps the rows a `=` filter of another consumer would have
pruned.
The skip is logged at DEBUG: it is a normal decision of the planner, not
an error.
### Release note
Fixed a bug where a query referencing a materialized CTE more than once
could silently return
fewer rows than expected (or no rows) when runtime filters were enabled,
for example when each
reference joins the CTE with an opposite range condition.
### Check List (For Author)
- Test: FE UT
(`RuntimeFilterTest#testPushSharedCteRuntimeFiltersWhichEveryConsumerApplies`,
`#testDoNotPushSharedCteRuntimeFiltersWhichOtherConsumersDoNotApply`,
`#testDoNotPushFiltersWithDifferentNullSemantics`, and the plan level
`#testPushSharedCteRuntimeFilterIntoTheProducer` /
`#testDoNotPushSingleConsumerCteRuntimeFilterIntoTheProducer`, which
assert on the runtime filters
installed inside the CTE producers of a planned query: the filter every
consumer applies reaches the
producer and is removed from the consumer it was built for, the filters
of a single consumer do not)
and the new regression test
`nereids_rules_p0/cte/test_cte_shared_producer_min_max_runtime_filter`
(four ordered cases plus a
NULL/right-outer case which mixes `<=>` and `=` on the same shared CTE),
which was run against an
unpatched FE (fails: `Check tag
'shared_cte_min_max_rf_opposite_directions' failed`) and against
the patched FE (passes). Manual test on a local cluster with internal
OLAP tables confirmed that
`runtime_filter_mode=GLOBAL/LOCAL` with `MIN_MAX` and a materialized CTE
now returns the same rows
as `OFF`. TPC-DS q95 was compared with and without this change
(`runtime_filter_type=12`, same
data and settings): the runtime filters on the inner scan of the CTE
body are unchanged and the
result is identical.
- Behavior changed: Yes. Only runtime filters that every consumer of a
materialized CTE applies are
pushed into the shared producer; such queries now return correct results
instead of empty ones.
- Does this need documentation: No
---
.../glue/translator/RuntimeFilterTranslator.java | 6 +-
.../processor/post/RuntimeFilterGenerator.java | 427 ++++++++++++++-------
.../post/RuntimeFilterPushDownVisitor.java | 33 +-
.../RuntimeFilterTranslatorBucketPruneTest.java | 30 ++
.../nereids/postprocess/RuntimeFilterTest.java | 364 +++++++++++++++++-
..._cte_shared_producer_min_max_runtime_filter.out | 33 ++
...e_shared_producer_min_max_runtime_filter.groovy | 148 +++++++
7 files changed, 885 insertions(+), 156 deletions(-)
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/RuntimeFilterTranslator.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/RuntimeFilterTranslator.java
index 8912502d1f4..f1ef04201d1 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/RuntimeFilterTranslator.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/RuntimeFilterTranslator.java
@@ -241,7 +241,11 @@ public class RuntimeFilterTranslator {
setPruningMetadata(origFilter, scanNode, group.get(i));
}
origFilter.setBloomFilterSizeCalculatedByNdv(head.isBloomFilterSizeCalculatedByNdv());
- setWaitTimeMs(origFilter, head.isNonBlocking(), isLocalTarget);
+ // The merged legacy filter applies on every target of the
group, so it may not wait when any
+ // of the filters it merges must not: the wait would put back
the edge that the non-blocking
+ // filter of the group removes, and a target of that filter
can then wait for a builder which
+ // itself waits for the target.
+ setWaitTimeMs(origFilter,
group.stream().anyMatch(RuntimeFilter::isNonBlocking), isLocalTarget);
org.apache.doris.planner.RuntimeFilter finalizedFilter =
finalize(origFilter);
scanNodeList.stream().filter(CTEScanNode.class::isInstance)
.forEach(f -> {
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/processor/post/RuntimeFilterGenerator.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/processor/post/RuntimeFilterGenerator.java
index 925bef6f760..2364d342f6a 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/processor/post/RuntimeFilterGenerator.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/processor/post/RuntimeFilterGenerator.java
@@ -28,6 +28,7 @@ import org.apache.doris.nereids.trees.expressions.GreaterThan;
import org.apache.doris.nereids.trees.expressions.GreaterThanEqual;
import org.apache.doris.nereids.trees.expressions.LessThan;
import org.apache.doris.nereids.trees.expressions.LessThanEqual;
+import org.apache.doris.nereids.trees.expressions.NullSafeEqual;
import org.apache.doris.nereids.trees.expressions.Slot;
import org.apache.doris.nereids.trees.expressions.SlotReference;
import org.apache.doris.nereids.trees.plans.AbstractPlan;
@@ -36,6 +37,7 @@ import org.apache.doris.nereids.trees.plans.Plan;
import org.apache.doris.nereids.trees.plans.algebra.Join;
import org.apache.doris.nereids.trees.plans.algebra.SetOperation;
import org.apache.doris.nereids.trees.plans.physical.AbstractPhysicalJoin;
+import org.apache.doris.nereids.trees.plans.physical.AbstractPhysicalPlan;
import org.apache.doris.nereids.trees.plans.physical.PhysicalCTEConsumer;
import org.apache.doris.nereids.trees.plans.physical.PhysicalCTEProducer;
import org.apache.doris.nereids.trees.plans.physical.PhysicalDistribute;
@@ -72,8 +74,10 @@ import org.apache.logging.log4j.Logger;
import java.util.HashMap;
import java.util.HashSet;
+import java.util.Iterator;
import java.util.List;
import java.util.Map;
+import java.util.Objects;
import java.util.Optional;
import java.util.Set;
import java.util.stream.Collectors;
@@ -111,136 +115,162 @@ public class RuntimeFilterGenerator extends
PlanPostProcessor {
}
Plan result = plan.accept(this, ctx);
- // try to push rf inside CTEProducer
- // collect cteProducers
+ pushRuntimeFiltersIntoCTEProducer(plan, ctx);
+ return result;
+ }
+
+ /**
+ * Push the runtime filters of the consumers of a CTE into its producer,
where one filter takes the
+ * place of the identical filters of all the consumers. See {@link
#selectPushableRuntimeFilters}.
+ */
+ private void pushRuntimeFiltersIntoCTEProducer(Plan plan, CascadesContext
ctx) {
RuntimeFilterContext rfCtx = ctx.getRuntimeFilterContext();
Map<CTEId, PhysicalCTEProducer> cteProducerMap =
plan.collect(PhysicalCTEProducer.class::isInstance)
.stream().collect(Collectors.toMap(p -> ((PhysicalCTEProducer)
p).getCteId(),
p -> (PhysicalCTEProducer) p));
- // collect cteConsumers which are RF targets
+ // collect the cte consumers that are runtime filter targets, grouped
by the cte they read
Map<CTEId, Set<PhysicalCTEConsumer>> cteIdToConsumersWithRF =
Maps.newHashMap();
Map<PhysicalCTEConsumer, Set<RuntimeFilter>> consumerToRFs =
Maps.newHashMap();
- Map<PhysicalCTEConsumer, Set<Expression>> consumerToSrcExpression =
Maps.newHashMap();
- List<RuntimeFilter> allRFs = rfCtx.getNereidsRuntimeFilter();
- for (RuntimeFilter rf : allRFs) {
- PhysicalRelation rel = rf.getTargetScan();
- if (rel instanceof PhysicalCTEConsumer) {
- PhysicalCTEConsumer consumer = (PhysicalCTEConsumer) rel;
- CTEId cteId = consumer.getCteId();
- cteIdToConsumersWithRF.computeIfAbsent(cteId, key ->
Sets.newHashSet()).add(consumer);
+ for (RuntimeFilter rf : rfCtx.getNereidsRuntimeFilter()) {
+ PhysicalRelation target = rf.getTargetScan();
+ if (target instanceof PhysicalCTEConsumer) {
+ PhysicalCTEConsumer consumer = (PhysicalCTEConsumer) target;
+ cteIdToConsumersWithRF.computeIfAbsent(consumer.getCteId(),
key -> Sets.newHashSet()).add(consumer);
consumerToRFs.computeIfAbsent(consumer, key ->
Sets.newHashSet()).add(rf);
- consumerToSrcExpression.computeIfAbsent(consumer, key ->
Sets.newHashSet())
- .add(rf.getSrcExpr());
- }
- }
- for (CTEId cteId : cteIdToConsumersWithRF.keySet()) {
- // if any consumer does not have RF, RF cannot be pushed down.
- // cteIdToConsumersWithRF.get(cteId).size() can not be 1, o.w.
this cte will be inlined.
- if (ctx.getCteIdToConsumers().get(cteId).size() ==
cteIdToConsumersWithRF.get(cteId).size()
- && cteIdToConsumersWithRF.get(cteId).size() >= 2) {
- // check if there is a common srcExpr among all the consumers
- Set<PhysicalCTEConsumer> consumers =
cteIdToConsumersWithRF.get(cteId);
- PhysicalCTEConsumer consumer0 = consumers.iterator().next();
- Set<Expression> candidateSrcExpressions =
consumerToSrcExpression.get(consumer0);
- for (PhysicalCTEConsumer currentConsumer : consumers) {
- Set<Expression> srcExpressionsOnCurrentConsumer =
consumerToSrcExpression.get(currentConsumer);
-
candidateSrcExpressions.retainAll(srcExpressionsOnCurrentConsumer);
- if (candidateSrcExpressions.isEmpty()) {
- break;
- }
- }
- if (!candidateSrcExpressions.isEmpty()) {
- // find RFs to push down
- for (Expression srcExpr : candidateSrcExpressions) {
- List<RuntimeFilter> rfsToPushDown =
Lists.newArrayList();
- for (PhysicalCTEConsumer consumer :
cteIdToConsumersWithRF.get(cteId)) {
- for (RuntimeFilter rf :
consumerToRFs.get(consumer)) {
- if (rf.getSrcExpr().equals(srcExpr)) {
- rfsToPushDown.add(rf);
- }
- }
- }
- if (rfsToPushDown.isEmpty()) {
- break;
- }
- if
(!canPushDownRuntimeFiltersIntoCTEProducer(rfsToPushDown, cteId)) {
- continue;
- }
+ }
+ }
+ for (Map.Entry<CTEId, Set<PhysicalCTEConsumer>> cteAndConsumers :
cteIdToConsumersWithRF.entrySet()) {
+ pushRuntimeFiltersIntoCTEProducer(cteAndConsumers.getKey(),
cteAndConsumers.getValue(),
+ consumerToRFs, ctx, rfCtx,
cteProducerMap.get(cteAndConsumers.getKey()));
+ }
+ }
- // the most right deep buildNode from rfsToPushDown is
used as buildNode for pushDown rf
- // since the srcExpr are the same, all buildNodes of
rfToPushDown are in the same tree path
- // the longest ancestors means its corresponding rf
build node is the most right deep one.
- List<RuntimeFilter> rightDeepRfs =
Lists.newArrayList();
- List<Plan> rightDeepAncestors =
rfsToPushDown.get(0).getBuilderNode().getAncestors();
- int rightDeepAncestorsSize = rightDeepAncestors.size();
- RuntimeFilter leftTop = rfsToPushDown.get(0);
- int leftTopAncestorsSize = rightDeepAncestorsSize;
- for (RuntimeFilter rf : rfsToPushDown) {
- List<Plan> ancestors =
rf.getBuilderNode().getAncestors();
- int currentAncestorsSize = ancestors.size();
- if (currentAncestorsSize >=
rightDeepAncestorsSize) {
- if (currentAncestorsSize ==
rightDeepAncestorsSize) {
- rightDeepRfs.add(rf);
- } else {
- rightDeepAncestorsSize =
currentAncestorsSize;
- rightDeepAncestors = ancestors;
- rightDeepRfs.clear();
- rightDeepRfs.add(rf);
- }
- }
- if (currentAncestorsSize < leftTopAncestorsSize) {
- leftTopAncestorsSize = currentAncestorsSize;
- leftTop = rf;
- }
- }
-
Preconditions.checkArgument(rightDeepAncestors.contains(leftTop.getBuilderNode()));
- // check nodes between right deep and left top are SPJ
and not denied join and not mark join
- boolean valid = true;
- for (Plan cursor : rightDeepAncestors) {
- if (cursor.equals(leftTop.getBuilderNode())) {
- break;
- }
- // valid = valid &&
SPJ_PLAN.contains(cursor.getClass());
- if (cursor instanceof AbstractPhysicalJoin) {
- AbstractPhysicalJoin cursorJoin =
(AbstractPhysicalJoin) cursor;
- valid =
(!RuntimeFilterGenerator.DENIED_JOIN_TYPES
- .contains(cursorJoin.getJoinType())
- || cursorJoin.isMarkJoin()) && valid;
- }
- if (!valid) {
- break;
- }
- }
+ /**
+ * Push the runtime filters of the consumers of one CTE into the producer
of that CTE.
+ */
+ private void pushRuntimeFiltersIntoCTEProducer(CTEId cteId,
Set<PhysicalCTEConsumer> consumers,
+ Map<PhysicalCTEConsumer, Set<RuntimeFilter>> consumerToRFs,
CascadesContext ctx,
+ RuntimeFilterContext rfCtx, PhysicalCTEProducer cteProducer) {
+ // if any consumer of this cte does not have a runtime filter, none of
them can be pushed down.
+ // there are always at least two consumers, otherwise this cte would
have been inlined.
+ if (consumers.size() < 2 ||
ctx.getCteIdToConsumers().get(cteId).size() != consumers.size()) {
+ return;
+ }
+ for (Expression srcExpr : commonSrcExpressions(consumers,
consumerToRFs)) {
+ List<RuntimeFilter> rfsOfSrcExpr =
runtimeFiltersOfSrcExpression(consumers, consumerToRFs, srcExpr);
+ for (List<RuntimeFilter> rfsOfIdentity :
selectPushableRuntimeFilters(rfsOfSrcExpr, consumers, cteId)) {
+ pushDownIdenticalFilters(rfsOfIdentity, cteId, rfCtx,
cteProducer);
+ }
+ }
+ }
- if (!valid) {
- break;
- }
+ /**
+ * The source expressions that every one of the given consumers has a
runtime filter for.
+ */
+ private static Set<Expression>
commonSrcExpressions(Set<PhysicalCTEConsumer> consumers,
+ Map<PhysicalCTEConsumer, Set<RuntimeFilter>> consumerToRFs) {
+ Iterator<PhysicalCTEConsumer> iterator = consumers.iterator();
+ Set<Expression> commonSrcExpressions =
srcExpressionsOf(consumerToRFs.get(iterator.next()));
+ while (iterator.hasNext() && !commonSrcExpressions.isEmpty()) {
+
commonSrcExpressions.retainAll(srcExpressionsOf(consumerToRFs.get(iterator.next())));
+ }
+ return commonSrcExpressions;
+ }
- for (RuntimeFilter rfToPush : rightDeepRfs) {
- Expression rightDeepTargetExpressionOnCTE = null;
- PhysicalRelation rel = rfToPush.getTargetScan();
- if (rel instanceof PhysicalCTEConsumer
- && ((PhysicalCTEConsumer)
rel).getCteId().equals(cteId)) {
- rightDeepTargetExpressionOnCTE =
rfToPush.getTargetExpression();
- }
-
- boolean pushedDown =
doPushDownIntoCTEProducerInternal(
- rfToPush,
- rightDeepTargetExpressionOnCTE,
- rfCtx,
- cteProducerMap.get(cteId)
- );
- if (pushedDown) {
- rfCtx.removeFilter(
- rfToPush,
-
rightDeepTargetExpressionOnCTE.getInputSlotExprIds().iterator().next());
- }
- }
- }
+ private static Set<Expression> srcExpressionsOf(Set<RuntimeFilter> rfs) {
+ return
rfs.stream().map(RuntimeFilter::getSrcExpr).collect(Collectors.toSet());
+ }
+
+ /**
+ * The runtime filters that all the given consumers have for one source
expression.
+ */
+ private static List<RuntimeFilter>
runtimeFiltersOfSrcExpression(Set<PhysicalCTEConsumer> consumers,
+ Map<PhysicalCTEConsumer, Set<RuntimeFilter>> consumerToRFs,
Expression srcExpr) {
+ List<RuntimeFilter> rfsOfSrcExpr = Lists.newArrayList();
+ for (PhysicalCTEConsumer consumer : consumers) {
+ for (RuntimeFilter rf : consumerToRFs.get(consumer)) {
+ if (rf.getSrcExpr().equals(srcExpr)) {
+ rfsOfSrcExpr.add(rf);
}
}
}
- return result;
+ Preconditions.checkArgument(!rfsOfSrcExpr.isEmpty());
+ return rfsOfSrcExpr;
+ }
+
+ /**
+ * Push one group of identical runtime filters into the shared CTE
producer. Only called with a group
+ * that every consumer applies, see {@link #selectPushableRuntimeFilters}.
+ */
+ private void pushDownIdenticalFilters(List<RuntimeFilter> rfsOfIdentity,
CTEId cteId,
+ RuntimeFilterContext rfCtx, PhysicalCTEProducer cteProducer) {
+ // the most right deep buildNode from rfsOfIdentity is used as
buildNode for pushDown rf
+ // since the srcExpr are the same, all buildNodes of rfsOfIdentity are
in the same tree path
+ // the longest ancestors means its corresponding rf build node is the
most right deep one.
+ List<RuntimeFilter> rightDeepRfs = Lists.newArrayList();
+ List<Plan> rightDeepAncestors =
rfsOfIdentity.get(0).getBuilderNode().getAncestors();
+ int rightDeepAncestorsSize = rightDeepAncestors.size();
+ RuntimeFilter leftTop = rfsOfIdentity.get(0);
+ int leftTopAncestorsSize = rightDeepAncestorsSize;
+ for (RuntimeFilter rf : rfsOfIdentity) {
+ List<Plan> ancestors = rf.getBuilderNode().getAncestors();
+ int currentAncestorsSize = ancestors.size();
+ if (currentAncestorsSize >= rightDeepAncestorsSize) {
+ if (currentAncestorsSize == rightDeepAncestorsSize) {
+ rightDeepRfs.add(rf);
+ } else {
+ rightDeepAncestorsSize = currentAncestorsSize;
+ rightDeepAncestors = ancestors;
+ rightDeepRfs.clear();
+ rightDeepRfs.add(rf);
+ }
+ }
+ if (currentAncestorsSize < leftTopAncestorsSize) {
+ leftTopAncestorsSize = currentAncestorsSize;
+ leftTop = rf;
+ }
+ }
+
Preconditions.checkArgument(rightDeepAncestors.contains(leftTop.getBuilderNode()));
+ // The filter of the deepest builder stands in for the filters of the
other consumers, which is sound
+ // only when it prunes at most as many rows as each of them would. The
source expression therefore has
+ // to keep the values it has on the build side of the deepest builder
while it travels up to the
+ // shallowest one: a node which can add a value to it -- the NULL an
outer join generates for the
+ // missing side of the source child, or the NULL a repeat synthesizes
for a grouping set which does
+ // not group by the source -- would make the filter below it prune the
rows the consumers above it
+ // still need.
+ if (!keepsSourceValue(rightDeepAncestors, leftTop.getBuilderNode())) {
+ return;
+ }
+ // The filter created on the producer stands in for every filter of
this group, so it has to keep the
+ // group's requirement not to be waited for. The producer feeds all
the consumers, while a filter which
+ // waits is produced by the build side of one of them: a replacement
which waits for a consumer whose
+ // own filter was non-blocking closes the cycle producer -> consumer
build side -> consumer -> producer,
+ // and the query stalls until the runtime filter or the query times
out. Waiting less than a member of
+ // the group asks for only makes the filter arrive later, it never
prunes a row the member would keep,
+ // so the replacement is non-blocking as soon as one member of the
group is.
+ boolean nonBlocking =
rfsOfIdentity.stream().anyMatch(RuntimeFilter::isNonBlocking);
+
+ for (RuntimeFilter rfToPush : rightDeepRfs) {
+ Expression rightDeepTargetExpressionOnCTE = null;
+ PhysicalRelation rel = rfToPush.getTargetScan();
+ if (rel instanceof PhysicalCTEConsumer
+ && ((PhysicalCTEConsumer) rel).getCteId().equals(cteId)) {
+ rightDeepTargetExpressionOnCTE =
rfToPush.getTargetExpression();
+ }
+
+ boolean pushedDown = doPushDownIntoCTEProducerInternal(
+ rfToPush,
+ rightDeepTargetExpressionOnCTE,
+ rfCtx,
+ cteProducer,
+ nonBlocking
+ );
+ if (pushedDown) {
+ rfCtx.removeFilter(
+ rfToPush,
+
rightDeepTargetExpressionOnCTE.getInputSlotExprIds().iterator().next());
+ }
+ }
}
/**
@@ -881,20 +911,153 @@ public class RuntimeFilterGenerator extends
PlanPostProcessor {
}
/**
- * Check whether runtime filters on CTE consumers can be pushed into their
shared CTE producer.
+ * Whether the nodes between the deepest builder and the given shallowest
builder only restrict the rows
+ * the source expression of the runtime filters is evaluated on, see
{@link #pushDownIdenticalFilters}.
+ * The list of ancestors starts with the deepest builder itself, whose
filter is the one which is pushed,
+ * so it is not part of the path.
+ */
+ @VisibleForTesting
+ public static boolean keepsSourceValue(List<Plan> ancestors, Plan
shallowestBuilder) {
+ if (ancestors.get(0).equals(shallowestBuilder)) {
+ // both filters are computed by the same node, their source values
are identical
+ return true;
+ }
+ for (int i = 1; i < ancestors.size(); i++) {
+ Plan node = ancestors.get(i);
+ if (node.equals(shallowestBuilder)) {
+ break;
+ }
+ if (SPJ_PLAN.stream().noneMatch(clazz -> clazz.isInstance(node))) {
+ return false;
+ }
+ if (node instanceof AbstractPhysicalJoin) {
+ AbstractPhysicalJoin<?, ?> join = (AbstractPhysicalJoin<?, ?>)
node;
+ if
(RuntimeFilterGenerator.DENIED_JOIN_TYPES.contains(join.getJoinType()) ||
join.isMarkJoin()) {
+ return false;
+ }
+ // the source is null extended when it comes from the side the
join generates NULLs for
+ if (isNullGeneratingChild(join.getJoinType(), join.child(0) ==
ancestors.get(i - 1))) {
+ return false;
+ }
+ }
+ }
+ return true;
+ }
+
+ /**
+ * Whether the given child of a join of that type is the side the join
generates NULLs for. Mirrors
+ * {@link RuntimeFilterPushDownVisitor}; the outer join types which are
denied already do not get here.
+ */
+ private static boolean isNullGeneratingChild(JoinType joinType, boolean
isLeftChild) {
+ if (joinType.isFullOuterJoin()) {
+ return true;
+ }
+ if (isLeftChild) {
+ return joinType.isRightOuterJoin() ||
joinType.isAsofRightOuterJoin();
+ }
+ return joinType.isLeftOuterJoin() || joinType.isAsofLeftOuterJoin();
+ }
+
+ /**
+ * Select the runtime filters of one source expression that may be pushed
into the shared CTE producer.
+ *
+ * <p>The producer feeds every consumer, so a filter may only be applied
on the producer when all the
+ * consumers apply the very same filter; a filter that only holds for one
consumer would prune the rows
+ * that the other consumers still need. The filters are therefore grouped
by identity -- same type, same
+ * min/max direction, same NULL semantics and same target expression on
the producer -- and only a group
+ * that every consumer applies is selected. For example, with `t c1 where
c1.k > b.x` and
+ * `t c2 where c2.k < b.x` the consumers produce a MIN and a MAX filter
on the same producer column:
+ * neither group covers both
+ * consumers, so neither is pushed. When on the other hand every consumer
applies the same pair of
+ * filters, for example a MIN_MAX and an IN_OR_BLOOM filter of the same
column, both groups are selected
+ * and each of them is still pushed once on the producer.
*/
@VisibleForTesting
- public static boolean canPushDownRuntimeFiltersIntoCTEProducer(
- List<RuntimeFilter> rfsToPushDown, CTEId cteId) {
- if (rfsToPushDown.isEmpty()) {
- LOG.warn("Skip pushing runtime filters into CTE producer because
no runtime filters exist for cteId: {}",
- cteId);
+ public static List<List<RuntimeFilter>> selectPushableRuntimeFilters(
+ List<RuntimeFilter> rfsOfSrcExpr, Set<PhysicalCTEConsumer>
consumers, CTEId cteId) {
+ Map<FilterIdentity, List<RuntimeFilter>> rfsByIdentity =
Maps.newLinkedHashMap();
+ for (RuntimeFilter rf : rfsOfSrcExpr) {
+ rfsByIdentity.computeIfAbsent(new FilterIdentity(rf, cteId), key
-> Lists.newArrayList()).add(rf);
+ }
+ List<List<RuntimeFilter>> pushable = Lists.newArrayList();
+ for (List<RuntimeFilter> rfsOfIdentity : rfsByIdentity.values()) {
+ Set<PhysicalCTEConsumer> consumersApplying = rfsOfIdentity.stream()
+ .map(rf -> (PhysicalCTEConsumer) rf.getTargetScan())
+ .collect(Collectors.toSet());
+ if (consumersApplying.size() == consumers.size()) {
+ pushable.add(rfsOfIdentity);
+ } else {
+ LOG.debug("Skip pushing runtime filters into CTE producer
because only {} of {} consumers"
+ + " of cteId: {} apply the filter {}, while
all of them have to apply it",
+ consumersApplying.size(), consumers.size(), cteId,
rfsOfIdentity.get(0));
+ }
+ }
+ return pushable;
+ }
+
+ /**
+ * Identity of a runtime filter with respect to a CTE producer. Two
filters with the same identity apply
+ * the very same predicate on the producer, therefore applying one of them
once on the producer is
+ * equivalent to applying it on every consumer. The source expression is
the same for all the filters
+ * compared here, so it does not take part in the identity.
+ */
+ private static final class FilterIdentity {
+ private final TRuntimeFilterType type;
+ private final TMinMaxRuntimeFilterType minMaxType;
+ private final boolean nullAware;
+ private final Expression producerTargetExpression;
+
+ private FilterIdentity(RuntimeFilter rf, CTEId cteId) {
+ this.type = rf.getType();
+ this.minMaxType = rf.gettMinMaxType();
+ this.nullAware = isNullAware(rf);
+ this.producerTargetExpression = getProducerTargetExpression(rf,
cteId);
+ }
+
+ @Override
+ public boolean equals(Object obj) {
+ if (this == obj) {
+ return true;
+ }
+ if (!(obj instanceof FilterIdentity)) {
+ return false;
+ }
+ FilterIdentity that = (FilterIdentity) obj;
+ return type == that.type && minMaxType == that.minMaxType &&
nullAware == that.nullAware
+ &&
producerTargetExpression.equals(that.producerTargetExpression);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(type, minMaxType, nullAware,
producerTargetExpression);
+ }
+
+ @Override
+ public String toString() {
+ return type + "/" + minMaxType + (nullAware ? "/nullAware" : "") +
" -> " + producerTargetExpression;
+ }
+ }
+
+ /**
+ * Whether the filter also keeps the rows whose probe column is NULL. The
legacy translation derives the
+ * mode from the builder of the filter -- a hash join conjunct with
EQ_FOR_NULL, or a set operation --
+ * see {@link org.apache.doris.planner.RuntimeFilter}, so it is a property
of the filter and not of the
+ * consumer which happens to push it down. Two filters which differ in it
must not be substituted for
+ * each other: a filter built from '=' prunes the rows of the NULL value
of its source, while the rows a
+ * '<=>' predicate matches are exactly those.
+ */
+ private static boolean isNullAware(RuntimeFilter rf) {
+ AbstractPhysicalPlan builder = rf.getBuilderNode();
+ if (builder instanceof PhysicalSetOperation) {
+ return true;
+ }
+ if (!(builder instanceof PhysicalHashJoin) || rf.getExprOrder() < 0) {
return false;
}
- Set<Expression> producerTargetExpressions = rfsToPushDown.stream()
- .map(rf -> getProducerTargetExpression(rf, cteId))
- .collect(Collectors.toSet());
- return producerTargetExpressions.size() == 1;
+ List<Expression> hashJoinConjuncts = ((PhysicalHashJoin<?, ?>)
builder).getHashJoinConjuncts();
+ Preconditions.checkArgument(rf.getExprOrder() <
hashJoinConjuncts.size(),
+ "exprOrder %s of the runtime filter is not a hash join
conjunct", rf.getExprOrder());
+ return hashJoinConjuncts.get(rf.getExprOrder()) instanceof
NullSafeEqual;
}
private static Expression getProducerTargetExpression(RuntimeFilter rf,
CTEId cteId) {
@@ -910,7 +1073,8 @@ public class RuntimeFilterGenerator extends
PlanPostProcessor {
}
private boolean doPushDownIntoCTEProducerInternal(RuntimeFilter rf,
Expression targetExpression,
- RuntimeFilterContext ctx,
PhysicalCTEProducer cteProducer) {
+ RuntimeFilterContext ctx,
PhysicalCTEProducer cteProducer,
+ boolean nonBlocking) {
PhysicalPlan inputPlanNode = (PhysicalPlan) cteProducer.child(0);
Slot unwrappedSlot = checkTargetChild(targetExpression);
if (unwrappedSlot == null) {
@@ -940,12 +1104,15 @@ public class RuntimeFilterGenerator extends
PlanPostProcessor {
if (!checkCanPushDownIntoBasicTable(inputPlanNode)) {
return false;
}
- // Use the PushDownVisitor to push inside the CTE producer subtree
+ // Use the PushDownVisitor to push inside the CTE producer subtree.
The non-blocking requirement of the
+ // filters this one replaces travels with it: the visitor creates a
new filter object, which would
+ // otherwise wait by default.
RuntimeFilterPushDownVisitor.PushDownContext pushDownContext =
RuntimeFilterPushDownVisitor.PushDownContext.createPushDownContext(
ctx, rf.getBuilderNode(), rf.getSrcExpr(),
producerTargetExpression,
rf.getType(), rf.gettMinMaxType(),
- !rf.isBloomFilterSizeCalculatedByNdv(),
rf.getBuildSideNdv(), rf.getExprOrder());
+ !rf.isBloomFilterSizeCalculatedByNdv(),
rf.getBuildSideNdv(), rf.getExprOrder())
+ .withNonBlocking(nonBlocking);
if (pushDownContext.isValid()) {
return inputPlanNode.accept(new RuntimeFilterPushDownVisitor(),
pushDownContext);
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/processor/post/RuntimeFilterPushDownVisitor.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/processor/post/RuntimeFilterPushDownVisitor.java
index 7f11df8aab0..7ea169dede3 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/processor/post/RuntimeFilterPushDownVisitor.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/processor/post/RuntimeFilterPushDownVisitor.java
@@ -72,11 +72,12 @@ public class RuntimeFilterPushDownVisitor extends
PlanVisitor<Boolean, PushDownC
final boolean hasUnknownColStats;
final long buildSideNdv;
final int exprOrder;
+ final boolean nonBlocking;
private PushDownContext(RuntimeFilterContext rfContext,
AbstractPhysicalPlan builderNode, Expression srcExpr,
Expression probeExpr,
TRuntimeFilterType type, TMinMaxRuntimeFilterType
singleSideMinMax,
- boolean hasUnknownColStats, long buildSideNdv, int exprOrder) {
+ boolean hasUnknownColStats, long buildSideNdv, int exprOrder,
boolean nonBlocking) {
this.rfContext = rfContext;
this.builderNode = builderNode;
this.srcExpr = srcExpr;
@@ -86,6 +87,7 @@ public class RuntimeFilterPushDownVisitor extends
PlanVisitor<Boolean, PushDownC
this.hasUnknownColStats = hasUnknownColStats;
this.buildSideNdv = buildSideNdv;
this.exprOrder = exprOrder;
+ this.nonBlocking = nonBlocking;
}
public static PushDownContext
createPushDownContext(RuntimeFilterContext rfContext,
@@ -112,7 +114,7 @@ public class RuntimeFilterPushDownVisitor extends
PlanVisitor<Boolean, PushDownC
TRuntimeFilterType type, TMinMaxRuntimeFilterType
singleSideMinMax,
boolean hasUnknownColStats, long buildSideNdv, int exprOrder) {
return new PushDownContext(rfContext, builderNode, srcExpr,
probeExpr,
- type, singleSideMinMax, hasUnknownColStats, buildSideNdv,
exprOrder);
+ type, singleSideMinMax, hasUnknownColStats, buildSideNdv,
exprOrder, false);
}
/**
@@ -125,7 +127,17 @@ public class RuntimeFilterPushDownVisitor extends
PlanVisitor<Boolean, PushDownC
public PushDownContext withNewProbeExpression(Expression newProbe) {
return new PushDownContext(rfContext, builderNode, srcExpr,
newProbe,
- type, singleSideMinMax, hasUnknownColStats, buildSideNdv,
exprOrder);
+ type, singleSideMinMax, hasUnknownColStats, buildSideNdv,
exprOrder, nonBlocking);
+ }
+
+ /**
+ * Carry the requirement not to be waited for into the filter which is
created for this context. The
+ * caller sets it when the created filter replaces filters which must
not be waited for, see
+ * {@code RuntimeFilterGenerator#pushDownIdenticalFilters}.
+ */
+ public PushDownContext withNonBlocking(boolean nonBlocking) {
+ return new PushDownContext(rfContext, builderNode, srcExpr,
probeExpr,
+ type, singleSideMinMax, hasUnknownColStats, buildSideNdv,
exprOrder, nonBlocking);
}
}
@@ -183,13 +195,19 @@ public class RuntimeFilterPushDownVisitor extends
PlanVisitor<Boolean, PushDownC
// V2-style: always create a separate RF per target.
// Dedup: skip if this scan already has an RF from the same (src,
type, builder).
- boolean alreadyApplied = scan.getAppliedRuntimeFilters().stream()
- .anyMatch(rf -> rf.getSrcExpr().equals(ctx.srcExpr)
+ RuntimeFilter alreadyApplied = scan.getAppliedRuntimeFilters().stream()
+ .filter(rf -> rf.getSrcExpr().equals(ctx.srcExpr)
&& rf.getType() == type
&& rf.getBuilderNode().equals(ctx.builderNode)
&& rf.getExprOrder() == ctx.exprOrder
- && rf.gettMinMaxType() == ctx.singleSideMinMax);
- if (alreadyApplied) {
+ && rf.gettMinMaxType() == ctx.singleSideMinMax)
+ .findFirst().orElse(null);
+ if (alreadyApplied != null) {
+ // The filter which is already applied on this scan stands in for
the one this context describes,
+ // so it has to keep the non-blocking requirement of both.
+ if (ctx.nonBlocking) {
+ alreadyApplied.setNonBlocking(true);
+ }
return true;
}
@@ -197,6 +215,7 @@ public class RuntimeFilterPushDownVisitor extends
PlanVisitor<Boolean, PushDownC
ctx.srcExpr, scanSlot, ctx.probeExpr,
type, ctx.exprOrder, ctx.builderNode, ctx.buildSideNdv,
!ctx.hasUnknownColStats, ctx.singleSideMinMax, scan);
+ filter.setNonBlocking(ctx.nonBlocking);
ctx.rfContext.generateRuntimeFilterPruneMetadata(filter);
scan.addAppliedRuntimeFilter(filter);
ctx.rfContext.addJoinToTargetMap(ctx.builderNode,
scanSlot.getExprId());
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/nereids/glue/translator/RuntimeFilterTranslatorBucketPruneTest.java
b/fe/fe-core/src/test/java/org/apache/doris/nereids/glue/translator/RuntimeFilterTranslatorBucketPruneTest.java
index 43cd6d7bae7..c83c732ebe3 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/nereids/glue/translator/RuntimeFilterTranslatorBucketPruneTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/nereids/glue/translator/RuntimeFilterTranslatorBucketPruneTest.java
@@ -150,6 +150,36 @@ class RuntimeFilterTranslatorBucketPruneTest {
Assertions.assertFalse(desc.isSetBucketPruningTargetIds());
}
+ /**
+ * One merged legacy filter applies on every target of its group, so it
may not wait when one of the
+ * filters it merges must not be waited for: the wait would come back to
the target of the non-blocking
+ * filter, which is exactly the wait edge that filter removes.
+ */
+ @Test
+ void testGroupedFiltersDoNotWaitWhenOneOfThemIsNonBlocking() {
+ // the blocking filter is the head of the group
+ Assertions.assertEquals(0,
translateGroupWithOneNonBlockingFilter(false).getWaitTimeMs());
+ // and it is not, so the group is not judged by the flag the head
happens to carry
+ Assertions.assertEquals(0,
translateGroupWithOneNonBlockingFilter(true).getWaitTimeMs());
+ }
+
+ /** Translate a group of two filters which only merge into one legacy
filter: one of them non-blocking. */
+ private TRuntimeFilterDesc translateGroupWithOneNonBlockingFilter(boolean
nonBlockingIsFirst) {
+ TranslatorHarness harness = new TranslatorHarness();
+ SlotReference blockingTarget = harness.addTargetSlot("dist_col",
harness.distributionColumn,
+ IntegerType.INSTANCE);
+ Column valueColumn = new Column("value_col", PrimitiveType.INT);
+ SlotReference nonBlockingTarget = harness.addTargetSlot("value_col",
valueColumn, IntegerType.INSTANCE);
+
+ RuntimeFilter blocking = harness.newFilter(blockingTarget,
blockingTarget);
+ RuntimeFilter nonBlocking = harness.newFilter(nonBlockingTarget,
nonBlockingTarget);
+ nonBlocking.setNonBlocking(true);
+
+ return harness.translate(nonBlockingIsFirst
+ ? ImmutableList.of(nonBlocking, blocking)
+ : ImmutableList.of(blocking, nonBlocking));
+ }
+
private static int firstLegacySlotId(TranslatorHarness harness,
SlotReference target) {
return
harness.translatorContext.findSlotRef(target.getExprId()).getSlotId().asInt();
}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/nereids/postprocess/RuntimeFilterTest.java
b/fe/fe-core/src/test/java/org/apache/doris/nereids/postprocess/RuntimeFilterTest.java
index 82f3ecc6a44..488c498c731 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/nereids/postprocess/RuntimeFilterTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/nereids/postprocess/RuntimeFilterTest.java
@@ -37,6 +37,7 @@ import org.apache.doris.nereids.trees.expressions.EqualTo;
import org.apache.doris.nereids.trees.expressions.ExprId;
import org.apache.doris.nereids.trees.expressions.Expression;
import org.apache.doris.nereids.trees.expressions.NamedExpression;
+import org.apache.doris.nereids.trees.expressions.NullSafeEqual;
import org.apache.doris.nereids.trees.expressions.Slot;
import org.apache.doris.nereids.trees.expressions.SlotReference;
import org.apache.doris.nereids.trees.expressions.Subtract;
@@ -49,10 +50,13 @@ import
org.apache.doris.nereids.trees.plans.commands.ExplainCommand;
import org.apache.doris.nereids.trees.plans.logical.LogicalPlan;
import org.apache.doris.nereids.trees.plans.physical.AbstractPhysicalPlan;
import org.apache.doris.nereids.trees.plans.physical.PhysicalCTEConsumer;
+import org.apache.doris.nereids.trees.plans.physical.PhysicalCTEProducer;
import org.apache.doris.nereids.trees.plans.physical.PhysicalHashJoin;
import org.apache.doris.nereids.trees.plans.physical.PhysicalOlapScan;
import org.apache.doris.nereids.trees.plans.physical.PhysicalPlan;
import org.apache.doris.nereids.trees.plans.physical.PhysicalProject;
+import org.apache.doris.nereids.trees.plans.physical.PhysicalRelation;
+import org.apache.doris.nereids.trees.plans.physical.PhysicalRepeat;
import org.apache.doris.nereids.trees.plans.physical.PhysicalSetOperation;
import org.apache.doris.nereids.trees.plans.physical.RuntimeFilter;
import org.apache.doris.nereids.types.IntegerType;
@@ -66,6 +70,7 @@ import org.apache.doris.thrift.TMinMaxRuntimeFilterType;
import org.apache.doris.thrift.TRuntimeFilterType;
import com.google.common.collect.ImmutableList;
+import com.google.common.collect.ImmutableSet;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
import org.mockito.Mockito;
@@ -73,6 +78,7 @@ import org.mockito.Mockito;
import java.util.ArrayList;
import java.util.List;
import java.util.Optional;
+import java.util.Set;
import java.util.function.Consumer;
import java.util.stream.Collectors;
@@ -780,37 +786,359 @@ public class RuntimeFilterTest extends SSBTestBase {
}
@Test
- public void
testPushSharedCteRuntimeFilterOnlyForSameProducerTargetExpression() {
+ public void testPushSharedCteRuntimeFiltersWhichEveryConsumerApplies() {
CTEId cteId = new CTEId(1);
SlotReference src = new SlotReference("src", IntegerType.INSTANCE);
SlotReference producerPk = new SlotReference("pk",
IntegerType.INSTANCE);
SlotReference consumerPk1 = new SlotReference("pk",
IntegerType.INSTANCE);
SlotReference consumerPk2 = new SlotReference("pk",
IntegerType.INSTANCE);
+ PhysicalCTEConsumer consumer1 = newCteConsumer(cteId, consumerPk1,
producerPk);
+ PhysicalCTEConsumer consumer2 = newCteConsumer(cteId, consumerPk2,
producerPk);
+ Set<PhysicalCTEConsumer> consumers = ImmutableSet.of(consumer1,
consumer2);
+
+ // Both consumers apply the same filter: it can be applied once on the
producer.
+ List<RuntimeFilter> sameFilter = ImmutableList.of(
+ newCteConsumerRuntimeFilter(consumer1, src, consumerPk1,
consumerPk1),
+ newCteConsumerRuntimeFilter(consumer2, src, consumerPk2,
consumerPk2));
+ Assertions.assertEquals(1,
RuntimeFilterGenerator.selectPushableRuntimeFilters(
+ sameFilter, consumers, cteId).size());
+
+ // Both consumers apply the same pair of filters, of two different
types: each filter is applied
+ // by every consumer, so both are still pushed, each as its own group.
+ List<RuntimeFilter> sameFilterPair = ImmutableList.of(
+ newCteConsumerRuntimeFilter(consumer1, src, consumerPk1,
consumerPk1,
+ TRuntimeFilterType.MIN_MAX,
TMinMaxRuntimeFilterType.MIN_MAX),
+ newCteConsumerRuntimeFilter(consumer1, src, consumerPk1,
consumerPk1,
+ TRuntimeFilterType.IN_OR_BLOOM,
TMinMaxRuntimeFilterType.MIN_MAX),
+ newCteConsumerRuntimeFilter(consumer2, src, consumerPk2,
consumerPk2,
+ TRuntimeFilterType.MIN_MAX,
TMinMaxRuntimeFilterType.MIN_MAX),
+ newCteConsumerRuntimeFilter(consumer2, src, consumerPk2,
consumerPk2,
+ TRuntimeFilterType.IN_OR_BLOOM,
TMinMaxRuntimeFilterType.MIN_MAX));
+ Assertions.assertEquals(2,
RuntimeFilterGenerator.selectPushableRuntimeFilters(
+ sameFilterPair, consumers, cteId).size());
+
+ // Only the filter that both consumers apply is pushed.
+ List<RuntimeFilter> partiallySharedFilter = ImmutableList.of(
+ newCteConsumerRuntimeFilter(consumer1, src, consumerPk1,
consumerPk1,
+ TRuntimeFilterType.MIN_MAX,
TMinMaxRuntimeFilterType.MIN_MAX),
+ newCteConsumerRuntimeFilter(consumer1, src, consumerPk1,
consumerPk1,
+ TRuntimeFilterType.IN_OR_BLOOM,
TMinMaxRuntimeFilterType.MIN_MAX),
+ newCteConsumerRuntimeFilter(consumer2, src, consumerPk2,
consumerPk2,
+ TRuntimeFilterType.IN_OR_BLOOM,
TMinMaxRuntimeFilterType.MIN_MAX));
+ List<List<RuntimeFilter>> pushable =
RuntimeFilterGenerator.selectPushableRuntimeFilters(
+ partiallySharedFilter, consumers, cteId);
+ Assertions.assertEquals(1, pushable.size());
+ Assertions.assertEquals(2, pushable.get(0).size());
+
+ // The filters target different expressions on the producer, so they
are not the same filter and
+ // neither of them is applied by both consumers.
+ List<RuntimeFilter> differentTargets = ImmutableList.of(
+ newCteConsumerRuntimeFilter(consumer1, src, consumerPk1,
+ new Add(consumerPk1, new IntegerLiteral(6))),
+ newCteConsumerRuntimeFilter(consumer2, src, consumerPk2,
+ new Subtract(consumerPk2, new IntegerLiteral(1))));
+
Assertions.assertTrue(RuntimeFilterGenerator.selectPushableRuntimeFilters(
+ differentTargets, consumers, cteId).isEmpty());
+ }
+
+ @Test
+ public void
testDoNotPushSharedCteRuntimeFiltersWhichOtherConsumersDoNotApply() {
+ CTEId cteId = new CTEId(1);
+ SlotReference src = new SlotReference("src", IntegerType.INSTANCE);
+ SlotReference producerK = new SlotReference("k", IntegerType.INSTANCE);
+ SlotReference consumerK1 = new SlotReference("k",
IntegerType.INSTANCE);
+ SlotReference consumerK2 = new SlotReference("k",
IntegerType.INSTANCE);
+ PhysicalCTEConsumer consumer1 = newCteConsumer(cteId, consumerK1,
producerK);
+ PhysicalCTEConsumer consumer2 = newCteConsumer(cteId, consumerK2,
producerK);
+ Set<PhysicalCTEConsumer> consumers = ImmutableSet.of(consumer1,
consumer2);
+
+ // `t c1 where c1.k > b.x` produces a MIN filter, `t c2 where c2.k <
b.x` produces a MAX filter.
+ // They target the same producer column and share the same source
expression, but each of them is
+ // applied by one consumer only: pushing them into the shared producer
would prune the rows that
+ // the other consumer still needs.
+ List<RuntimeFilter> oppositeMinMaxFilters = ImmutableList.of(
+ newCteConsumerRuntimeFilter(consumer1, src, consumerK1,
consumerK1,
+ TRuntimeFilterType.MIN_MAX,
TMinMaxRuntimeFilterType.MIN),
+ newCteConsumerRuntimeFilter(consumer2, src, consumerK2,
consumerK2,
+ TRuntimeFilterType.MIN_MAX,
TMinMaxRuntimeFilterType.MAX));
+
Assertions.assertTrue(RuntimeFilterGenerator.selectPushableRuntimeFilters(
+ oppositeMinMaxFilters, consumers, cteId).isEmpty());
+
+ // Same direction: both consumers apply the same filter, so it is
pushed.
+ List<RuntimeFilter> sameDirectionFilters = ImmutableList.of(
+ newCteConsumerRuntimeFilter(consumer1, src, consumerK1,
consumerK1,
+ TRuntimeFilterType.MIN_MAX,
TMinMaxRuntimeFilterType.MIN),
+ newCteConsumerRuntimeFilter(consumer2, src, consumerK2,
consumerK2,
+ TRuntimeFilterType.MIN_MAX,
TMinMaxRuntimeFilterType.MIN));
+ Assertions.assertEquals(1,
RuntimeFilterGenerator.selectPushableRuntimeFilters(
+ sameDirectionFilters, consumers, cteId).size());
+
+ // Different filter types, each applied by one consumer only: neither
is pushed.
+ List<RuntimeFilter> differentTypeFilters = ImmutableList.of(
+ newCteConsumerRuntimeFilter(consumer1, src, consumerK1,
consumerK1,
+ TRuntimeFilterType.MIN_MAX,
TMinMaxRuntimeFilterType.MIN),
+ newCteConsumerRuntimeFilter(consumer2, src, consumerK2,
consumerK2,
+ TRuntimeFilterType.IN_OR_BLOOM,
TMinMaxRuntimeFilterType.MIN));
+
Assertions.assertTrue(RuntimeFilterGenerator.selectPushableRuntimeFilters(
+ differentTypeFilters, consumers, cteId).isEmpty());
+ }
+
+ @Test
+ public void testDoNotPushFiltersWithDifferentNullSemantics() {
+ CTEId cteId = new CTEId(1);
+ SlotReference src = new SlotReference("src", IntegerType.INSTANCE);
+ SlotReference producerK = new SlotReference("k", IntegerType.INSTANCE);
+ SlotReference consumerK1 = new SlotReference("k",
IntegerType.INSTANCE);
+ SlotReference consumerK2 = new SlotReference("k",
IntegerType.INSTANCE);
+ PhysicalCTEConsumer consumer1 = newCteConsumer(cteId, consumerK1,
producerK);
+ PhysicalCTEConsumer consumer2 = newCteConsumer(cteId, consumerK2,
producerK);
+ Set<PhysicalCTEConsumer> consumers = ImmutableSet.of(consumer1,
consumer2);
+
+ // `c1.k <=> b.x` produces a null aware filter, `c2.k = b.x` an
ordinary one. They prune different
+ // rows -- the ordinary one removes the rows whose probe column is
NULL, and those are exactly the
+ // rows the null aware predicate matches -- so neither may replace the
other on the shared producer.
+ List<RuntimeFilter> differentNullSemantics = ImmutableList.of(
+ newCteConsumerRuntimeFilter(consumer1, src, consumerK1,
consumerK1, newHashJoinBuilder(true)),
+ newCteConsumerRuntimeFilter(consumer2, src, consumerK2,
consumerK2, newHashJoinBuilder(false)));
+
Assertions.assertTrue(RuntimeFilterGenerator.selectPushableRuntimeFilters(
+ differentNullSemantics, consumers, cteId).isEmpty());
+
+ // The same NULL semantics on every consumer: the filter is pushed.
+ List<RuntimeFilter> sameNullSemantics = ImmutableList.of(
+ newCteConsumerRuntimeFilter(consumer1, src, consumerK1,
consumerK1, newHashJoinBuilder(false)),
+ newCteConsumerRuntimeFilter(consumer2, src, consumerK2,
consumerK2, newHashJoinBuilder(false)));
+ Assertions.assertEquals(1,
RuntimeFilterGenerator.selectPushableRuntimeFilters(
+ sameNullSemantics, consumers, cteId).size());
+ }
+
+ /**
+ * The runtime filters which every consumer of a CTE applies are pushed
into the producer, where one of
+ * them filters the rows of all the consumers; the filters of a single
consumer must stay where they are,
+ * otherwise they prune the rows the other consumers still need.
+ */
+ @Test
+ public void testPushSharedCteRuntimeFilterIntoTheProducer() {
+ int oldType =
connectContext.getSessionVariable().getRuntimeFilterType();
+ boolean oldMaterialize =
connectContext.getSessionVariable().enableCTEMaterialize;
+
connectContext.getSessionVariable().setRuntimeFilterType(TRuntimeFilterType.MIN_MAX.getValue());
+ connectContext.getSessionVariable().enableCTEMaterialize = true;
+ try {
+ // Both consumers apply the same MIN filter built from
`p_partkey`, so it reaches the producer.
+ PhysicalPlan plan = planAfterPostProcess(
+ "with t as (select lo_partkey as k from lineorder)"
+ + " select c2.k from t c2 cross join t c1 cross
join part"
+ + " where c1.k > p_partkey and c2.k > p_partkey");
+ List<RuntimeFilter> pushedIntoProducer =
runtimeFiltersInsideCteProducers(plan);
+ Assertions.assertFalse(pushedIntoProducer.isEmpty(),
+ "the filter which every consumer applies must be pushed
into the producer");
+ // the filter which was moved into the producer is no longer
applied by the consumer it was built
+ // for, otherwise the rows it filters would be filtered twice
+ List<RuntimeFilter> onConsumers =
runtimeFiltersOnCteConsumers(plan);
+
Assertions.assertTrue(pushedIntoProducer.stream().noneMatch(onConsumers::contains),
+ () -> "a filter pushed into the producer must not stay on
its consumer: " + onConsumers);
+ } finally {
+ connectContext.getSessionVariable().setRuntimeFilterType(oldType);
+ connectContext.getSessionVariable().enableCTEMaterialize =
oldMaterialize;
+ }
+ }
+
+ @Test
+ public void testDoNotPushSingleConsumerCteRuntimeFilterIntoTheProducer() {
+ int oldType =
connectContext.getSessionVariable().getRuntimeFilterType();
+ boolean oldMaterialize =
connectContext.getSessionVariable().enableCTEMaterialize;
+
connectContext.getSessionVariable().setRuntimeFilterType(TRuntimeFilterType.MIN_MAX.getValue());
+ connectContext.getSessionVariable().enableCTEMaterialize = true;
+ try {
+ // `c1.k > p_partkey` builds a MIN filter and `c2.k < p_partkey` a
MAX one: each of them is
+ // applied by one consumer only, so neither may be applied on the
shared producer.
+ List<RuntimeFilter> pushedIntoProducer =
runtimeFiltersInsideCteProducers(planAfterPostProcess(
+ "with t as (select lo_partkey as k from lineorder)"
+ + " select c2.k from t c2 cross join t c1 cross
join part"
+ + " where c1.k > p_partkey and c2.k < p_partkey"));
+ Assertions.assertTrue(pushedIntoProducer.isEmpty(),
+ () -> "the filters of a single consumer must not reach the
producer: " + pushedIntoProducer);
+ } finally {
+ connectContext.getSessionVariable().setRuntimeFilterType(oldType);
+ connectContext.getSessionVariable().enableCTEMaterialize =
oldMaterialize;
+ }
+ }
+
+ /**
+ * The filter created on the shared producer replaces the filters of the
consumers, so it has to keep the
+ * requirement of that group not to be waited for. The producer feeds
every consumer, while a filter which
+ * waits is built by one of them: a replacement which waits for a consumer
whose own filter was
+ * non-blocking closes the cycle producer -> build side of a consumer ->
consumer -> producer, and the
+ * query stalls until the runtime filter or the query times out.
+ */
+ @Test
+ public void
testPushSharedCteRuntimeFilterIntoTheProducerKeepsTheNonBlockingRequirement() {
+ boolean oldMaterialize =
connectContext.getSessionVariable().enableCTEMaterialize;
+ boolean oldExpandByInnerJoin =
connectContext.getSessionVariable().expandRuntimeFilterByInnerJoin;
+ boolean oldDecoupled =
connectContext.getSessionVariable().enableDecoupledRuntimeFilter;
+ long oldMinDecoupledRows =
connectContext.getSessionVariable().minDecoupledRfTargetRows;
+ connectContext.getSessionVariable().enableCTEMaterialize = true;
+ connectContext.getSessionVariable().expandRuntimeFilterByInnerJoin =
true;
+ connectContext.getSessionVariable().enableDecoupledRuntimeFilter =
true;
+ // the tables of the test catalog hold one row, while the decoupled
filter of the plan below targets a
+ // scan of it: keep the filter which describes the behavior under test
rather than the pruning of
+ // filters which cannot arrive in time on a tiny scan.
+ connectContext.getSessionVariable().minDecoupledRfTargetRows = 0;
+ try {
+ // `c1.c_custkey = s_suppkey` is the condition join: its standard
filter targets the CTE consumer
+ // `c1` and expands to `c2`, while the reverse decoupled filter is
built by the deeper join
+ // `c2.c_custkey = c1.c_custkey`, whose build side carries the
filter on `c_region`. The decoupled
+ // filter is therefore preferred and both standard filters are
marked non-blocking.
+ PhysicalPlan plan = planAfterPostProcess(
+ "with t as (select c_custkey, c_region from customer)"
+ + " select c1.c_custkey from t c2 join t c1 on
c2.c_custkey = c1.c_custkey"
+ + " join supplier on c1.c_custkey = s_suppkey"
+ + " where c1.c_region = 'ASIA'");
+ List<RuntimeFilter> pushedIntoProducer =
runtimeFiltersInsideCteProducers(plan);
+ Assertions.assertFalse(pushedIntoProducer.isEmpty(),
+ () -> "the filter which every consumer applies must be
pushed into the producer: "
+ + plan.treeString());
+
Assertions.assertTrue(pushedIntoProducer.stream().allMatch(RuntimeFilter::isNonBlocking),
+ () -> "a filter pushed into the producer must keep the
non-blocking requirement of the"
+ + " filters it replaces: " + pushedIntoProducer);
+ } finally {
+ connectContext.getSessionVariable().enableCTEMaterialize =
oldMaterialize;
+ connectContext.getSessionVariable().expandRuntimeFilterByInnerJoin
= oldExpandByInnerJoin;
+ connectContext.getSessionVariable().enableDecoupledRuntimeFilter =
oldDecoupled;
+ connectContext.getSessionVariable().minDecoupledRfTargetRows =
oldMinDecoupledRows;
+ }
+ }
+
+ /**
+ * A node which synthesizes a value of the source column -- here the NULL
the repeat adds for the grouping
+ * set which does not group by it -- makes a filter built below it prune
the rows the consumers above it
+ * still need, so the deepest filter must not stand in for them.
+ */
+ @Test
+ public void testDoNotPushRuntimeFilterWhichCrossesAValueSynthesizingNode()
{
+ boolean oldMaterialize =
connectContext.getSessionVariable().enableCTEMaterialize;
+ connectContext.getSessionVariable().enableCTEMaterialize = true;
+ try {
+ PhysicalPlan plan = planAfterPostProcess(
+ "with t as (select lo_partkey as k from lineorder)"
+ + " select c1.k from t c1 join ("
+ + " select x from ("
+ + " select c2.k as k2, p.p_partkey as x from t
c2 join part p on c2.k <=> p.p_partkey"
+ + " ) a group by grouping sets ((x), ())"
+ + " ) g on c1.k <=> g.x");
+ Assertions.assertTrue(plan.containsType(PhysicalRepeat.class),
+ "the query must plan the grouping sets which synthesize
the NULL");
+
Assertions.assertTrue(!plan.<Plan>collect(PhysicalCTEProducer.class::isInstance).isEmpty(),
+ "the query must materialize the CTE");
+
Assertions.assertFalse(runtimeFiltersOnCteConsumers(plan).isEmpty(),
+ "the consumers must keep the filters which may not be
pushed");
+ List<RuntimeFilter> pushedIntoProducer =
runtimeFiltersInsideCteProducers(plan);
+ Assertions.assertTrue(pushedIntoProducer.isEmpty(),
+ () -> "a filter which crosses a value synthesizing node
must not reach the producer: "
+ + pushedIntoProducer);
+ } finally {
+ connectContext.getSessionVariable().enableCTEMaterialize =
oldMaterialize;
+ }
+ }
+
+ /**
+ * A right outer join null extends its left child, so a filter whose
source comes from that side did not
+ * observe the NULL the join adds: it may not stand in for a filter which
is built above the join, whose
+ * build side does contain that NULL. The same join keeps the source
values when the source comes from the
+ * right child it preserves, and an inner join never adds a value.
+ */
+ @Test
+ public void testSourceValueIsNotKeptWhenAnOuterJoinNullExtendsIt() {
+ Plan deepestBuilder = Mockito.mock(Plan.class);
+ Plan shallowestBuilder = Mockito.mock(Plan.class);
+ PhysicalHashJoin<?, ?> rightOuterJoin =
newMockJoin(JoinType.RIGHT_OUTER_JOIN, deepestBuilder);
+ // the source comes from the left child, which the right outer join
null extends
+ Assertions.assertFalse(RuntimeFilterGenerator.keepsSourceValue(
+ ImmutableList.of(deepestBuilder, rightOuterJoin),
shallowestBuilder));
+ // the same join, with the source coming from the right child it
preserves
+ Assertions.assertTrue(RuntimeFilterGenerator.keepsSourceValue(
+ ImmutableList.of(rightOuterJoin.child(1), rightOuterJoin),
shallowestBuilder));
+ // an inner join restricts the rows of the source, it never adds a
value to it
+ Assertions.assertTrue(RuntimeFilterGenerator.keepsSourceValue(
+ ImmutableList.of(deepestBuilder,
newMockJoin(JoinType.INNER_JOIN, deepestBuilder)),
+ shallowestBuilder));
+ }
+
+ /** A join node whose left child is the given plan. */
+ private PhysicalHashJoin<?, ?> newMockJoin(JoinType joinType, Plan
leftChild) {
+ PhysicalHashJoin<?, ?> join = Mockito.mock(PhysicalHashJoin.class);
+ Mockito.when(join.getJoinType()).thenReturn(joinType);
+ Mockito.when(join.child(0)).thenReturn(leftChild);
+ Mockito.when(join.child(1)).thenReturn(Mockito.mock(Plan.class));
+ return join;
+ }
- List<RuntimeFilter> sameTargetFilters = ImmutableList.of(
- newCteConsumerRuntimeFilter(src, consumerPk1, consumerPk1,
producerPk, cteId),
- newCteConsumerRuntimeFilter(src, consumerPk2, consumerPk2,
producerPk, cteId));
-
Assertions.assertTrue(RuntimeFilterGenerator.canPushDownRuntimeFiltersIntoCTEProducer(
- sameTargetFilters, cteId));
+ private PhysicalPlan planAfterPostProcess(String sql) {
+ PlanChecker checker =
PlanChecker.from(connectContext).analyze(sql).rewrite().optimize();
+ return new
PlanPostProcessors(checker.getCascadesContext()).process(checker.getBestPlanTree());
+ }
+
+ /** The runtime filters which are still applied by the consumers of a CTE.
*/
+ private static List<RuntimeFilter>
runtimeFiltersOnCteConsumers(PhysicalPlan plan) {
+ List<RuntimeFilter> applied = new ArrayList<>();
+ for (Plan consumer :
plan.<Plan>collect(PhysicalCTEConsumer.class::isInstance)) {
+ applied.addAll(((AbstractPhysicalPlan)
consumer).getAppliedRuntimeFilters());
+ }
+ return applied;
+ }
+
+ /** The runtime filters which were installed on the relations inside the
CTE producers. */
+ private static List<RuntimeFilter>
runtimeFiltersInsideCteProducers(PhysicalPlan plan) {
+ List<RuntimeFilter> applied = new ArrayList<>();
+ for (Plan producer :
plan.<Plan>collect(PhysicalCTEProducer.class::isInstance)) {
+ for (Plan relation : ((PhysicalCTEProducer<?>) producer).child(0)
+ .<Plan>collect(PhysicalRelation.class::isInstance)) {
+ applied.addAll(((AbstractPhysicalPlan)
relation).getAppliedRuntimeFilters());
+ }
+ }
+ return applied;
+ }
- List<RuntimeFilter> differentTargetFilters = ImmutableList.of(
- newCteConsumerRuntimeFilter(src, consumerPk1,
- new Add(consumerPk1, new IntegerLiteral(6)),
producerPk, cteId),
- newCteConsumerRuntimeFilter(src, consumerPk2,
- new Subtract(consumerPk2, new IntegerLiteral(1)),
producerPk, cteId));
-
Assertions.assertFalse(RuntimeFilterGenerator.canPushDownRuntimeFiltersIntoCTEProducer(
- differentTargetFilters, cteId));
+ /** A hash join whose only conjunct is an EQ_FOR_NULL one when
nullSafeEqual is set. */
+ private AbstractPhysicalPlan newHashJoinBuilder(boolean nullSafeEqual) {
+ PhysicalHashJoin<?, ?> join = Mockito.mock(PhysicalHashJoin.class);
+ SlotReference left = new SlotReference("k", IntegerType.INSTANCE);
+ SlotReference right = new SlotReference("x", IntegerType.INSTANCE);
+ Mockito.when(join.getHashJoinConjuncts()).thenReturn(ImmutableList.of(
+ nullSafeEqual ? new NullSafeEqual(left, right) : new
EqualTo(left, right)));
+ return join;
}
- private RuntimeFilter newCteConsumerRuntimeFilter(Expression src, Slot
targetSlot,
- Expression targetExpression, Slot producerSlot, CTEId cteId) {
+ private PhysicalCTEConsumer newCteConsumer(CTEId cteId, Slot targetSlot,
Slot producerSlot) {
PhysicalCTEConsumer consumer = Mockito.mock(PhysicalCTEConsumer.class);
Mockito.when(consumer.getCteId()).thenReturn(cteId);
Mockito.when(consumer.getProducerSlot(targetSlot)).thenReturn(producerSlot);
- AbstractPhysicalPlan builder =
Mockito.mock(AbstractPhysicalPlan.class);
+ return consumer;
+ }
+
+ private RuntimeFilter newCteConsumerRuntimeFilter(PhysicalCTEConsumer
consumer, Expression src,
+ Slot targetSlot, Expression targetExpression) {
+ return newCteConsumerRuntimeFilter(consumer, src, targetSlot,
targetExpression,
+ TRuntimeFilterType.IN_OR_BLOOM,
TMinMaxRuntimeFilterType.MIN_MAX);
+ }
+
+ private RuntimeFilter newCteConsumerRuntimeFilter(PhysicalCTEConsumer
consumer, Expression src,
+ Slot targetSlot, Expression targetExpression, TRuntimeFilterType
type,
+ TMinMaxRuntimeFilterType minMaxType) {
+ return newCteConsumerRuntimeFilter(consumer, src, targetSlot,
targetExpression,
+ Mockito.mock(AbstractPhysicalPlan.class), type, minMaxType);
+ }
+
+ private RuntimeFilter newCteConsumerRuntimeFilter(PhysicalCTEConsumer
consumer, Expression src,
+ Slot targetSlot, Expression targetExpression, AbstractPhysicalPlan
builder) {
+ return newCteConsumerRuntimeFilter(consumer, src, targetSlot,
targetExpression,
+ builder, TRuntimeFilterType.IN_OR_BLOOM,
TMinMaxRuntimeFilterType.MIN_MAX);
+ }
+
+ private RuntimeFilter newCteConsumerRuntimeFilter(PhysicalCTEConsumer
consumer, Expression src,
+ Slot targetSlot, Expression targetExpression, AbstractPhysicalPlan
builder,
+ TRuntimeFilterType type, TMinMaxRuntimeFilterType minMaxType) {
return new
RuntimeFilter(RuntimeFilterId.createGenerator().getNextId(), src, targetSlot,
targetExpression,
- TRuntimeFilterType.IN_OR_BLOOM, 0, builder, -1L, true,
- TMinMaxRuntimeFilterType.MIN_MAX, consumer);
+ type, 0, builder, -1L, true, minMaxType, consumer);
}
@Test
diff --git
a/regression-test/data/nereids_rules_p0/cte/test_cte_shared_producer_min_max_runtime_filter.out
b/regression-test/data/nereids_rules_p0/cte/test_cte_shared_producer_min_max_runtime_filter.out
new file mode 100644
index 00000000000..a057df8a0a7
--- /dev/null
+++
b/regression-test/data/nereids_rules_p0/cte/test_cte_shared_producer_min_max_runtime_filter.out
@@ -0,0 +1,33 @@
+-- This file is automatically generated. You should know what you did if you
want to edit this
+-- !shared_cte_min_max_rf_opposite_directions --
+0 10 0 10 3
+0 5 0 5 3
+1 10 1 10 3
+1 5 1 5 3
+
+-- !shared_cte_min_max_rf_same_direction --
+10 10 10 10 3
+10 5 10 5 3
+5 10 5 10 3
+5 5 5 5 3
+
+-- !shared_cte_min_max_rf_inlined_cte --
+0 10 0 10 3
+0 5 0 5 3
+1 10 1 10 3
+1 5 1 5 3
+
+-- !shared_cte_null_aware --
+\N \N
+3 3
+
+-- !shared_cte_grouping_sets --
+\N \N
+3 3
+
+-- !shared_cte_min_max_rf_off --
+0 10 0 10 3
+0 5 0 5 3
+1 10 1 10 3
+1 5 1 5 3
+
diff --git
a/regression-test/suites/nereids_rules_p0/cte/test_cte_shared_producer_min_max_runtime_filter.groovy
b/regression-test/suites/nereids_rules_p0/cte/test_cte_shared_producer_min_max_runtime_filter.groovy
new file mode 100644
index 00000000000..4db23ec803f
--- /dev/null
+++
b/regression-test/suites/nereids_rules_p0/cte/test_cte_shared_producer_min_max_runtime_filter.groovy
@@ -0,0 +1,148 @@
+// 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.
+
+suite("test_cte_shared_producer_min_max_runtime_filter") {
+ sql "SET enable_nereids_planner=true"
+ sql "SET enable_fallback_to_original_planner=false"
+ sql "SET enable_pipeline_engine=true"
+ // Materialize the CTE so that both references share one producer scan.
+ sql "SET enable_cte_materialize=true"
+ sql "SET inline_cte_referenced_threshold=0"
+ sql "SET runtime_filter_type='MIN_MAX'"
+ sql "SET enable_runtime_filter_prune=false"
+ sql "SET runtime_filter_wait_time_ms=10000"
+
+ sql "DROP TABLE IF EXISTS cte_rf_shared_producer_f"
+ sql """
+ CREATE TABLE cte_rf_shared_producer_f (
+ k INT
+ ) ENGINE=OLAP
+ DUPLICATE KEY(k)
+ DISTRIBUTED BY HASH(k) BUCKETS 2
+ PROPERTIES ("replication_num" = "1")
+ """
+ sql "INSERT INTO cte_rf_shared_producer_f VALUES (0), (1), (5), (10)"
+
+ sql "DROP TABLE IF EXISTS cte_rf_shared_producer_b"
+ sql """
+ CREATE TABLE cte_rf_shared_producer_b (
+ x INT
+ ) ENGINE=OLAP
+ DUPLICATE KEY(x)
+ DISTRIBUTED BY HASH(x) BUCKETS 1
+ PROPERTIES ("replication_num" = "1")
+ """
+ sql "INSERT INTO cte_rf_shared_producer_b VALUES (3)"
+
+ sql "DROP TABLE IF EXISTS cte_rf_null_f"
+ sql """
+ CREATE TABLE cte_rf_null_f (
+ k INT
+ ) ENGINE=OLAP
+ DUPLICATE KEY(k)
+ DISTRIBUTED BY HASH(k) BUCKETS 2
+ PROPERTIES ("replication_num" = "1")
+ """
+ sql "INSERT INTO cte_rf_null_f VALUES (NULL), (1), (3)"
+
+ sql "DROP TABLE IF EXISTS cte_rf_null_b"
+ sql """
+ CREATE TABLE cte_rf_null_b (
+ x INT
+ ) ENGINE=OLAP
+ DUPLICATE KEY(x)
+ DISTRIBUTED BY HASH(x) BUCKETS 1
+ PROPERTIES ("replication_num" = "1")
+ """
+ sql "INSERT INTO cte_rf_null_b VALUES (NULL), (3)"
+
+ // The two references of the CTE need disjoint row ranges: `c1.k > b.x`
asks for a MIN runtime
+ // filter and `c2.k < b.x` for a MAX one. Both target the same column of
the shared producer, so
+ // pushing both of them into the producer prunes every row (k >= 3 AND k
<= 3) and the query
+ // silently loses the rows each consumer still needs.
+ sql "SET runtime_filter_mode='GLOBAL'"
+ order_qt_shared_cte_min_max_rf_opposite_directions """
+ WITH t AS (SELECT k, ABS(k) AS v FROM cte_rf_shared_producer_f)
+ SELECT c2.k AS lo, c1.k AS hi, c2.v AS lo_v, c1.v AS hi_v, b.x
+ FROM t c2 CROSS JOIN t c1 CROSS JOIN cte_rf_shared_producer_b b
+ WHERE c1.k > b.x AND c2.k < b.x
+ ORDER BY lo, hi
+ """
+
+ // Both references filter in the same direction, so both consumers produce
the same runtime
+ // filter. Pushing it into the producer is equivalent to applying it on
every consumer.
+ order_qt_shared_cte_min_max_rf_same_direction """
+ WITH t AS (SELECT k, ABS(k) AS v FROM cte_rf_shared_producer_f)
+ SELECT c2.k AS lo, c1.k AS hi, c2.v AS lo_v, c1.v AS hi_v, b.x
+ FROM t c2 CROSS JOIN t c1 CROSS JOIN cte_rf_shared_producer_b b
+ WHERE c1.k > b.x AND c2.k > b.x
+ ORDER BY lo, hi
+ """
+
+ // Without a shared producer each reference gets its own scan and its own
runtime filter.
+ sql "SET enable_cte_materialize=false"
+ order_qt_shared_cte_min_max_rf_inlined_cte """
+ WITH t AS (SELECT k, ABS(k) AS v FROM cte_rf_shared_producer_f)
+ SELECT c2.k AS lo, c1.k AS hi, c2.v AS lo_v, c1.v AS hi_v, b.x
+ FROM t c2 CROSS JOIN t c1 CROSS JOIN cte_rf_shared_producer_b b
+ WHERE c1.k > b.x AND c2.k < b.x
+ ORDER BY lo, hi
+ """
+ sql "SET enable_cte_materialize=true"
+
+ // NULL semantics are part of the filter identity: `c1.k <=> s.x` produces
a null aware filter while the
+ // deeper `c2.k = b0.x` produces an ordinary one. They prune different
rows -- the ordinary filter removes
+ // the rows whose probe column is NULL, which are exactly the rows the
null aware predicate matches -- so
+ // neither may be applied on the shared producer, otherwise the row which
matches NULL with NULL is lost.
+ sql "SET runtime_filter_type=12"
+ order_qt_shared_cte_null_aware """
+ WITH t AS (SELECT k FROM cte_rf_null_f)
+ SELECT c1.k AS a, s.x AS bx
+ FROM t c1
+ JOIN (
+ SELECT c2.k AS k2, b0.x AS x
+ FROM t c2 RIGHT OUTER JOIN cte_rf_null_b b0 ON c2.k = b0.x
+ ) s ON c1.k <=> s.x
+ ORDER BY a, bx
+ """
+
+ // A value synthesizing node between the two builders makes the deeper
filter prune rows the upper
+ // consumer still needs: the repeat adds the NULL of the grouping set
which does not group by x, so the
+ // filter built below it, from `b.x = {3}`, must not be applied on the
shared producer. The row which
+ // matches that synthesized NULL with the NULL of the CTE is lost when it
is.
+ order_qt_shared_cte_grouping_sets """
+ WITH t AS (SELECT k FROM cte_rf_null_f)
+ SELECT c1.k AS a, g.x AS gx
+ FROM t c1
+ JOIN (
+ SELECT x FROM (
+ SELECT c2.k AS k2, b.x AS x
+ FROM t c2 JOIN cte_rf_shared_producer_b b ON c2.k <=> b.x
+ ) z GROUP BY GROUPING SETS ((x), ())
+ ) g ON c1.k <=> g.x
+ ORDER BY a, gx
+ """
+
+ sql "SET runtime_filter_mode='OFF'"
+ order_qt_shared_cte_min_max_rf_off """
+ WITH t AS (SELECT k, ABS(k) AS v FROM cte_rf_shared_producer_f)
+ SELECT c2.k AS lo, c1.k AS hi, c2.v AS lo_v, c1.v AS hi_v, b.x
+ FROM t c2 CROSS JOIN t c1 CROSS JOIN cte_rf_shared_producer_b b
+ WHERE c1.k > b.x AND c2.k < b.x
+ ORDER BY lo, hi
+ """
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]