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]
