This is an automated email from the ASF dual-hosted git repository.

HappenLee pushed a commit to branch 4.1_performance
in repository https://gitbox.apache.org/repos/asf/doris.git

commit 70e1241fbaafcc640ba0513c4792a6718007ffb1
Author: happenlee <[email protected]>
AuthorDate: Mon Aug 17 12:03:16 2026 +0800

    [improvement](fe) Skip oversized Bloom filters for shuffle joins
    
    ### What problem does this PR solve?
    
    Issue Number: None
    
    Related PR: None
    
    Problem Summary: Shuffle joins can build and merge large Bloom runtime 
filters even when the build-side cardinality exceeds 
runtime_filter_max_build_row_count. Add an opt-in session variable that skips 
BLOOM and IN_OR_BLOOM filters for partitioned and bucket-shuffle joins when 
build statistics are unknown or above the threshold, while preserving 
broadcast, colocate, and MIN_MAX filters. A zero threshold continues to force 
generation.
    
    ### Release note
    
    Add enable_ignore_runtime_filter_for_large_shuffle_join to optionally skip 
oversized shuffle-join Bloom runtime filters.
    
    ### Check List (For Author)
    
    - Test: Unit Test
        - FE_UT_PARALLEL=48 ./run-fe-ut.sh --run 
org.apache.doris.nereids.postprocess.RuntimeFilterTest
        - ./build.sh --fe
    - Behavior changed: Yes. When enabled, oversized or unknown build-side 
shuffle joins no longer generate Bloom-class runtime filters.
    - Does this need documentation: No
---
 .../processor/post/RuntimeFilterGenerator.java     | 26 ++++++++++++++
 .../java/org/apache/doris/qe/SessionVariable.java  | 17 +++++++++
 .../nereids/postprocess/RuntimeFilterTest.java     | 41 ++++++++++++++++++++++
 3 files changed, 84 insertions(+)

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 bbc926bbd11..729d419e8fb 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
@@ -265,8 +265,10 @@ public class RuntimeFilterGenerator extends 
PlanPostProcessor {
             return join;
         }
         RuntimeFilterContext ctx = context.getRuntimeFilterContext();
+        boolean ignoreBloomFilter = 
shouldIgnoreBloomFilterForShuffleJoin(join, ctx);
         List<TRuntimeFilterType> legalTypes = 
RuntimeFilterTypeHelper.getSupportedRuntimeFilterTypes().stream()
                 .filter(type -> 
ctx.getSessionVariable().allowedRuntimeFilterType(type))
+                .filter(type -> !ignoreBloomFilter || !isBloomFilter(type))
                 .collect(Collectors.toList());
 
         List<Expression> hashJoinConjuncts = join.getHashJoinConjuncts();
@@ -296,6 +298,30 @@ public class RuntimeFilterGenerator extends 
PlanPostProcessor {
         return join;
     }
 
+    private boolean shouldIgnoreBloomFilterForShuffleJoin(
+            PhysicalHashJoin<? extends Plan, ? extends Plan> join, 
RuntimeFilterContext ctx) {
+        if 
(!ctx.getSessionVariable().isEnableIgnoreRuntimeFilterForLargeShuffleJoin()) {
+            return false;
+        }
+        Join.ShuffleType shuffleType = join.shuffleType();
+        if (shuffleType != Join.ShuffleType.shuffle
+                && shuffleType != Join.ShuffleType.bucketShuffle
+                && shuffleType != Join.ShuffleType.shuffleBucket) {
+            return false;
+        }
+        long maxBuildRowCount = 
ctx.getSessionVariable().runtimeFilterMaxBuildRowCount;
+        if (maxBuildRowCount <= 0) {
+            return false;
+        }
+        AbstractPlan buildSide = (AbstractPlan) join.right();
+        double buildRowCount = buildSide.getStats() == null ? -1 : 
buildSide.getStats().getRowCount();
+        return buildRowCount <= 0 || buildRowCount > maxBuildRowCount;
+    }
+
+    private boolean isBloomFilter(TRuntimeFilterType type) {
+        return type == TRuntimeFilterType.BLOOM || type == 
TRuntimeFilterType.IN_OR_BLOOM;
+    }
+
     /**
      *
      * T1 join T1 on T1.a=T2.a where T1.a=1 and T2.a=1
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/SessionVariable.java 
b/fe/fe-core/src/main/java/org/apache/doris/qe/SessionVariable.java
index 28b33167980..4a08179ec62 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/SessionVariable.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/SessionVariable.java
@@ -266,6 +266,8 @@ public class SessionVariable implements Serializable, 
Writable {
     public static final String ENABLE_SYNC_RUNTIME_FILTER_SIZE = 
"enable_sync_runtime_filter_size";
     public static final String RUNTIME_FILTER_TREE_PUBLISH_MAX_SEND_BYTES =
             "runtime_filter_tree_publish_max_send_bytes";
+    public static final String 
ENABLE_IGNORE_RUNTIME_FILTER_FOR_LARGE_SHUFFLE_JOIN =
+            "enable_ignore_runtime_filter_for_large_shuffle_join";
 
     public static final String ENABLE_PARALLEL_RESULT_SINK = 
"enable_parallel_result_sink";
 
@@ -1770,6 +1772,12 @@ public class SessionVariable implements Serializable, 
Writable {
     @VariableMgr.VarAttr(name = "runtime_filter_max_build_row_count", 
needForward = true, fuzzy = false)
     public long runtimeFilterMaxBuildRowCount = 64L * 1024L * 1024L;
 
+    @VariableMgr.VarAttr(name = 
ENABLE_IGNORE_RUNTIME_FILTER_FOR_LARGE_SHUFFLE_JOIN, needForward = true,
+            description = {"是否忽略 build 侧行数未知或超过 
runtime_filter_max_build_row_count 的 shuffle join Bloom Filter",
+                    "Whether to ignore Bloom filters for shuffle joins whose 
build-side row count is unknown or "
+                            + "exceeds runtime_filter_max_build_row_count"})
+    private boolean enableIgnoreRuntimeFilterForLargeShuffleJoin = false;
+
     @VariableMgr.VarAttr(name = ENABLE_PARALLEL_RESULT_SINK, needForward = 
true, fuzzy = true)
     private boolean enableParallelResultSink = true;
 
@@ -4797,6 +4805,15 @@ public class SessionVariable implements Serializable, 
Writable {
         return runtimeFilterTreePublishMaxSendBytes;
     }
 
+    public boolean isEnableIgnoreRuntimeFilterForLargeShuffleJoin() {
+        return enableIgnoreRuntimeFilterForLargeShuffleJoin;
+    }
+
+    public void setEnableIgnoreRuntimeFilterForLargeShuffleJoin(
+            boolean enableIgnoreRuntimeFilterForLargeShuffleJoin) {
+        this.enableIgnoreRuntimeFilterForLargeShuffleJoin = 
enableIgnoreRuntimeFilterForLargeShuffleJoin;
+    }
+
     public void setEnableLocalShuffle(boolean enableLocalShuffle) {
         this.enableLocalShuffle = enableLocalShuffle;
     }
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 59538f98e22..7b053a67f05 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
@@ -35,6 +35,7 @@ import org.apache.doris.nereids.trees.expressions.ExprId;
 import org.apache.doris.nereids.trees.expressions.NamedExpression;
 import org.apache.doris.nereids.trees.expressions.SlotReference;
 import org.apache.doris.nereids.trees.expressions.literal.NullLiteral;
+import org.apache.doris.nereids.trees.plans.AbstractPlan;
 import org.apache.doris.nereids.trees.plans.DistributeType;
 import org.apache.doris.nereids.trees.plans.JoinType;
 import org.apache.doris.nereids.trees.plans.Plan;
@@ -49,6 +50,7 @@ import 
org.apache.doris.nereids.trees.plans.physical.RuntimeFilter;
 import org.apache.doris.nereids.util.MemoTestUtils;
 import org.apache.doris.nereids.util.PlanChecker;
 import org.apache.doris.qe.OriginStatement;
+import org.apache.doris.qe.SessionVariable;
 
 import com.google.common.collect.ImmutableList;
 import com.google.common.collect.Sets;
@@ -59,6 +61,7 @@ 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;
 
 public class RuntimeFilterTest extends SSBTestBase {
@@ -82,6 +85,39 @@ public class RuntimeFilterTest extends SSBTestBase {
                 Pair.of("c_custkey", "lo_custkey")));
     }
 
+    @Test
+    public void testIgnoreRuntimeFilterForLargeShuffleJoin() {
+        SessionVariable sessionVariable = connectContext.getSessionVariable();
+        long originalMaxBuildRowCount = 
sessionVariable.runtimeFilterMaxBuildRowCount;
+        boolean originalEnableIgnore = 
sessionVariable.isEnableIgnoreRuntimeFilterForLargeShuffleJoin();
+        String shuffleSql = "SELECT * FROM lineorder JOIN [shuffle] customer"
+                + " ON c_custkey = lo_custkey";
+        String broadcastSql = "SELECT * FROM lineorder JOIN [broadcast] 
customer"
+                + " ON c_custkey = lo_custkey";
+        Consumer<PhysicalPlan> setLargeBuildSide = plan -> {
+            PhysicalHashJoin<?, ?> join = (PhysicalHashJoin<?, ?>) 
plan.collect(PhysicalHashJoin.class::isInstance)
+                    .iterator().next();
+            AbstractPlan buildSide = (AbstractPlan) join.right();
+            buildSide.setStatistics(buildSide.getStats().withRowCount(2));
+        };
+        try {
+            sessionVariable.runtimeFilterMaxBuildRowCount = 1;
+
+            
sessionVariable.setEnableIgnoreRuntimeFilterForLargeShuffleJoin(false);
+            Assertions.assertEquals(1, getRuntimeFilters(shuffleSql, 
setLargeBuildSide).get().size());
+
+            
sessionVariable.setEnableIgnoreRuntimeFilterForLargeShuffleJoin(true);
+            Assertions.assertEquals(0, getRuntimeFilters(shuffleSql, 
setLargeBuildSide).get().size());
+            Assertions.assertEquals(1, getRuntimeFilters(broadcastSql, 
setLargeBuildSide).get().size());
+
+            sessionVariable.runtimeFilterMaxBuildRowCount = 0;
+            Assertions.assertEquals(1, getRuntimeFilters(shuffleSql, 
setLargeBuildSide).get().size());
+        } finally {
+            sessionVariable.runtimeFilterMaxBuildRowCount = 
originalMaxBuildRowCount;
+            
sessionVariable.setEnableIgnoreRuntimeFilterForLargeShuffleJoin(originalEnableIgnore);
+        }
+    }
+
     @Test
     public void testGenerateRuntimeFilterByIllegalSrcExpr() {
         String sql = "SELECT * FROM lineorder JOIN customer on c_custkey = 
c_custkey";
@@ -333,11 +369,16 @@ public class RuntimeFilterTest extends SSBTestBase {
     }
 
     private Optional<List<RuntimeFilter>> getRuntimeFilters(String sql) {
+        return getRuntimeFilters(sql, plan -> { });
+    }
+
+    private Optional<List<RuntimeFilter>> getRuntimeFilters(String sql, 
Consumer<PhysicalPlan> beforePostProcess) {
         PlanChecker checker = PlanChecker.from(connectContext)
                 .analyze(sql)
                 .rewrite()
                 .optimize();
         PhysicalPlan plan = checker.getBestPlanTree();
+        beforePostProcess.accept(plan);
         plan = new 
PlanPostProcessors(checker.getCascadesContext()).process(plan);
         System.out.println(plan.treeString());
         new PhysicalPlanTranslator(new 
PlanTranslatorContext(checker.getCascadesContext())).translatePlan(plan);


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to