From 36272e7f65dbe9b13d010c645ede0a9404c81d0c Mon Sep 17 00:00:00 2001 From: lambxu Date: Sat, 29 Aug 2026 04:23:35 +0800 Subject: [PATCH 1/3] save --- .../doris/nereids/rules/rewrite/SkewJoin.java | 58 ++++++++++--------- .../org/apache/doris/qe/SessionVariable.java | 14 +++++ 2 files changed, 45 insertions(+), 27 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/SkewJoin.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/SkewJoin.java index 5db8f4217f54bb..8c9ecc89077e03 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/SkewJoin.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/SkewJoin.java @@ -70,7 +70,9 @@ private static Plan transform(MatchingContext> ctx) { LogicalJoin join = ctx.root; Expression skewExpr = null; List hotValues = new ArrayList<>(); - if (join.getHashJoinConjuncts().size() != 1) { + boolean allowMultiKey = + ConnectContext.get().getSessionVariable().isEnableSkewJoinMultiKey(); + if (!allowMultiKey && join.getHashJoinConjuncts().size() != 1) { return null; } AbstractPlan left = (AbstractPlan) join.left(); @@ -82,35 +84,37 @@ private static Plan transform(MatchingContext> ctx) { right.accept(derive, new DeriveContext()); } - EqualPredicate equal = (EqualPredicate) join.getHashJoinConjuncts().get(0); - if (join.left().getOutputSet().contains(equal.right())) { - equal = equal.commute(); - } - - if (join.getJoinType().isInnerJoin() || join.getJoinType().isLeftOuterJoin()) { - Expression leftEqHand = equal.child(0); - if (left.getStats().findColumnStatistics(leftEqHand) != null) { - ColumnStatistic leftColStats = left.getStats().findColumnStatistics(leftEqHand); - Map filtered = StatisticsUtil.getHotValuesWithOriginalThreshold( - leftColStats.getHotValues(), leftColStats.ndv); - if (filtered != null) { - skewExpr = leftEqHand; - hotValues.addAll(filtered.keySet()); - } + for (Expression conjunct : join.getHashJoinConjuncts()) { + EqualPredicate equal = (EqualPredicate) conjunct; + if (join.left().getOutputSet().contains(equal.right())) { + equal = equal.commute(); } - } else if (join.getJoinType().isRightOuterJoin()) { - Expression rightEqHand = equal.child(1); - if (right.getStats().findColumnStatistics(rightEqHand) != null) { - ColumnStatistic rightColStats = right.getStats().findColumnStatistics(rightEqHand); - Map filtered = StatisticsUtil.getHotValuesWithOriginalThreshold( - rightColStats.getHotValues(), rightColStats.ndv); - if (filtered != null) { - skewExpr = rightEqHand; - hotValues.addAll(filtered.keySet()); + Map filtered = null; + Expression sideExpr = null; + if (join.getJoinType().isInnerJoin() || join.getJoinType().isLeftOuterJoin()) { + Expression leftEqHand = equal.child(0); + if (left.getStats().findColumnStatistics(leftEqHand) != null) { + ColumnStatistic leftColStats = left.getStats().findColumnStatistics(leftEqHand); + filtered = StatisticsUtil.getHotValuesWithOriginalThreshold( + leftColStats.getHotValues(), leftColStats.ndv); + sideExpr = leftEqHand; + } + } else if (join.getJoinType().isRightOuterJoin()) { + Expression rightEqHand = equal.child(1); + if (right.getStats().findColumnStatistics(rightEqHand) != null) { + ColumnStatistic rightColStats = right.getStats().findColumnStatistics(rightEqHand); + filtered = StatisticsUtil.getHotValuesWithOriginalThreshold( + rightColStats.getHotValues(), rightColStats.ndv); + sideExpr = rightEqHand; } + } else { + return null; + } + if (filtered != null && !filtered.isEmpty()) { + skewExpr = sideExpr; + hotValues.addAll(filtered.keySet()); + break; } - } else { - return null; } if (skewExpr == null || hotValues.isEmpty()) { return null; 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 c5b503e9f144b0..2e2e8cfc855eeb 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 @@ -326,6 +326,9 @@ public String toString() { // turn off all automatic join reorder algorithms public static final String DISABLE_JOIN_REORDER = "disable_join_reorder"; + // Allow SkewJoin auto-salt to trigger on joins with more than one hash key. + public static final String ENABLE_SKEW_JOIN_MULTI_KEY = "enable_skew_join_multi_key"; + public static final String MAX_JOIN_NUMBER_OF_REORDER = "max_join_number_of_reorder"; public static final String ENABLE_NEREIDS_DML = "enable_nereids_dml"; @@ -1949,6 +1952,9 @@ public void setCboNetWeight(double cboNetWeight) { @VarAttrDef.VarAttr(name = DISABLE_JOIN_REORDER) private boolean disableJoinReorder = false; + @VariableMgr.VarAttr(name = ENABLE_SKEW_JOIN_MULTI_KEY, needForward = true) + private boolean enableSkewJoinMultiKey = false; + @VarAttrDef.VarAttr(name = MAX_JOIN_NUMBER_OF_REORDER) private int maxJoinNumberOfReorder = 63; @@ -5017,6 +5023,14 @@ public boolean isDisableJoinReorder() { return disableJoinReorder; } + public boolean isEnableSkewJoinMultiKey() { + return enableSkewJoinMultiKey; + } + + public void setEnableSkewJoinMultiKey(boolean enable) { + this.enableSkewJoinMultiKey = enable; + } + public boolean isEnableBushyTree() { return enableBushyTree; } From a9308138889e2a128c0f77692da8e648ae5114ac Mon Sep 17 00:00:00 2001 From: lambxu Date: Sat, 29 Aug 2026 04:26:38 +0800 Subject: [PATCH 2/3] save --- .../src/main/java/org/apache/doris/qe/SessionVariable.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 2e2e8cfc855eeb..df916f4cc53d18 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 @@ -1953,7 +1953,7 @@ public void setCboNetWeight(double cboNetWeight) { private boolean disableJoinReorder = false; @VariableMgr.VarAttr(name = ENABLE_SKEW_JOIN_MULTI_KEY, needForward = true) - private boolean enableSkewJoinMultiKey = false; + private boolean enableSkewJoinMultiKey = true; @VarAttrDef.VarAttr(name = MAX_JOIN_NUMBER_OF_REORDER) private int maxJoinNumberOfReorder = 63; From 1fbeaf96e50c72bcedff4500a27eabd3da087331 Mon Sep 17 00:00:00 2001 From: lambxu Date: Sat, 29 Aug 2026 07:04:10 +0800 Subject: [PATCH 3/3] save --- .../src/main/java/org/apache/doris/qe/SessionVariable.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 df916f4cc53d18..65441c5dd26d5a 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 @@ -1952,7 +1952,7 @@ public void setCboNetWeight(double cboNetWeight) { @VarAttrDef.VarAttr(name = DISABLE_JOIN_REORDER) private boolean disableJoinReorder = false; - @VariableMgr.VarAttr(name = ENABLE_SKEW_JOIN_MULTI_KEY, needForward = true) + @VarAttrDef.VarAttr(name = ENABLE_SKEW_JOIN_MULTI_KEY, needForward = true) private boolean enableSkewJoinMultiKey = true; @VarAttrDef.VarAttr(name = MAX_JOIN_NUMBER_OF_REORDER)