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 &gt; b.x` and
+     * `t c2 where c2.k &lt; 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
+     * '&lt;=&gt;' 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]

Reply via email to