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 ced019606794bf46d352b6213b74b0f0794f4a14
Author: happenlee <[email protected]>
AuthorDate: Mon Aug 17 12:19:17 2026 +0800

    [fix](fe) Exclude bucket shuffle joins from Bloom filter skipping
    
    ### What problem does this PR solve?
    
    Issue Number: None
    
    Related PR: None
    
    Problem Summary: The oversized Bloom runtime filter guard also covered 
bucket-shuffle joins. Restrict the guard to two-sided shuffle joins so 
bucketShuffle and shuffleBucket plans continue generating Bloom filters. Extend 
the unit test with a real bucket-shuffle plan and verify that its filter 
remains.
    
    ### Release note
    
    The large shuffle-join Bloom filter guard no longer applies to 
bucket-shuffle joins.
    
    ### 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. Bucket-shuffle joins retain Bloom runtime filters 
even when the large shuffle-join guard is enabled.
    - Does this need documentation: No
---
 .../nereids/processor/post/RuntimeFilterGenerator.java   |  5 +----
 .../main/java/org/apache/doris/qe/SessionVariable.java   |  6 +++---
 .../doris/nereids/postprocess/RuntimeFilterTest.java     | 16 ++++++++++++++++
 3 files changed, 20 insertions(+), 7 deletions(-)

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 729d419e8fb..40f7517d9e9 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
@@ -303,10 +303,7 @@ public class RuntimeFilterGenerator extends 
PlanPostProcessor {
         if 
(!ctx.getSessionVariable().isEnableIgnoreRuntimeFilterForLargeShuffleJoin()) {
             return false;
         }
-        Join.ShuffleType shuffleType = join.shuffleType();
-        if (shuffleType != Join.ShuffleType.shuffle
-                && shuffleType != Join.ShuffleType.bucketShuffle
-                && shuffleType != Join.ShuffleType.shuffleBucket) {
+        if (join.shuffleType() != Join.ShuffleType.shuffle) {
             return false;
         }
         long maxBuildRowCount = 
ctx.getSessionVariable().runtimeFilterMaxBuildRowCount;
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 4a08179ec62..8db225a1274 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
@@ -1773,9 +1773,9 @@ public class SessionVariable implements Serializable, 
Writable {
     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"})
+            description = {"是否忽略 build 侧行数未知或超过 
runtime_filter_max_build_row_count 的双边 shuffle join Bloom Filter",
+                    "Whether to ignore Bloom filters for two-sided 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)
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 7b053a67f05..d65e8c6bdaa 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
@@ -39,6 +39,7 @@ 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;
+import org.apache.doris.nereids.trees.plans.algebra.Join;
 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;
@@ -89,9 +90,12 @@ public class RuntimeFilterTest extends SSBTestBase {
     public void testIgnoreRuntimeFilterForLargeShuffleJoin() {
         SessionVariable sessionVariable = connectContext.getSessionVariable();
         long originalMaxBuildRowCount = 
sessionVariable.runtimeFilterMaxBuildRowCount;
+        double originalBroadcastRowCountLimit = 
sessionVariable.getBroadcastRowCountLimit();
         boolean originalEnableIgnore = 
sessionVariable.isEnableIgnoreRuntimeFilterForLargeShuffleJoin();
         String shuffleSql = "SELECT * FROM lineorder JOIN [shuffle] customer"
                 + " ON c_custkey = lo_custkey";
+        String bucketShuffleSql = "SELECT * FROM lineorder JOIN customer"
+                + " ON c_custkey = lo_orderkey";
         String broadcastSql = "SELECT * FROM lineorder JOIN [broadcast] 
customer"
                 + " ON c_custkey = lo_custkey";
         Consumer<PhysicalPlan> setLargeBuildSide = plan -> {
@@ -100,20 +104,32 @@ public class RuntimeFilterTest extends SSBTestBase {
             AbstractPlan buildSide = (AbstractPlan) join.right();
             buildSide.setStatistics(buildSide.getStats().withRowCount(2));
         };
+        Consumer<PhysicalPlan> setLargeBucketShuffleBuildSide = plan -> {
+            PhysicalHashJoin<?, ?> join = (PhysicalHashJoin<?, ?>) 
plan.collect(PhysicalHashJoin.class::isInstance)
+                    .iterator().next();
+            Assertions.assertTrue(join.shuffleType() == 
Join.ShuffleType.bucketShuffle
+                    || join.shuffleType() == Join.ShuffleType.shuffleBucket);
+            AbstractPlan buildSide = (AbstractPlan) join.right();
+            buildSide.setStatistics(buildSide.getStats().withRowCount(2));
+        };
         try {
             sessionVariable.runtimeFilterMaxBuildRowCount = 1;
+            sessionVariable.setBroadcastRowCountLimit(0);
 
             
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(bucketShuffleSql, 
setLargeBucketShuffleBuildSide).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.setBroadcastRowCountLimit(originalBroadcastRowCountLimit);
             
sessionVariable.setEnableIgnoreRuntimeFilterForLargeShuffleJoin(originalEnableIgnore);
         }
     }


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

Reply via email to