From 5ee9e8993cf27e16bcb5a31fe9521193c8b6616b Mon Sep 17 00:00:00 2001 From: yangtao555 Date: Tue, 4 Aug 2026 20:18:58 +0800 Subject: [PATCH] [fix](mv) Distinguish partition compensation from UNION ALL rewrite --- .../mv/AbstractMaterializedViewRule.java | 11 ++++++----- ...erializedViewAggregateOnNoneAggregateRule.java | 4 ++-- .../exploration/mv/PartitionCompensator.java | 7 +++---- .../exploration/mv/PartitionCompensatorTest.java | 15 +++++++++++++++ 4 files changed, 26 insertions(+), 11 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/exploration/mv/AbstractMaterializedViewRule.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/exploration/mv/AbstractMaterializedViewRule.java index 38c4fade7ff859..535f4eb9eaff97 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/exploration/mv/AbstractMaterializedViewRule.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/exploration/mv/AbstractMaterializedViewRule.java @@ -375,17 +375,18 @@ protected List doRewrite(StructInfo queryStructInfo, CascadesContext casca // if mv can not offer any partition for query, query rewrite bail out to avoid cycle run return rewriteResults; } - boolean partitionNeedUnion = PartitionCompensator.needUnionRewrite(invalidPartitions, cascadesContext); - boolean canUnionRewrite = canUnionRewrite(queryPlan, - (AsyncMaterializationContext) materializationContext, cascadesContext); - if (partitionNeedUnion && !canUnionRewrite) { + boolean hasPartitionCompensation = + PartitionCompensator.hasPartitionCompensation(invalidPartitions); + boolean needBaseTableUnion = !invalidPartitions.value().isEmpty(); + if (needBaseTableUnion && !canUnionRewrite(queryPlan, + (AsyncMaterializationContext) materializationContext, cascadesContext)) { materializationContext.recordFailReason(queryStructInfo, "need compensate union all, but can not, because the query structInfo", () -> String.format("mv partition info is %s, and the query plan is %s", mtmv.getMvPartitionInfo(), queryPlan.treeString())); return rewriteResults; } - if (partitionNeedUnion) { + if (hasPartitionCompensation) { Pair planAndNeedAddFilterPair = StructInfo.addFilterOnTableScan(queryPlan, invalidPartitions.value(), cascadesContext); if (planAndNeedAddFilterPair == null) { diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/exploration/mv/MaterializedViewAggregateOnNoneAggregateRule.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/exploration/mv/MaterializedViewAggregateOnNoneAggregateRule.java index 12a822b73885ee..a1fcd4c9d53e20 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/exploration/mv/MaterializedViewAggregateOnNoneAggregateRule.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/exploration/mv/MaterializedViewAggregateOnNoneAggregateRule.java @@ -105,8 +105,8 @@ protected Pair>, Map>> c Pair>, Map>> invalidPartitions = super.calcInvalidPartitions(queryUsedBaseTablePartitionMap, rewrittenPlan, cascadesContext, materializationContext); - if (PartitionCompensator.needUnionRewrite(invalidPartitions, cascadesContext)) { - // if query use some invalid partition in mv, bail out + if (PartitionCompensator.hasPartitionCompensation(invalidPartitions)) { + // Aggregate-on-non-aggregate rewrite does not support partition compensation. return null; } return invalidPartitions; diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/exploration/mv/PartitionCompensator.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/exploration/mv/PartitionCompensator.java index fe2e88cd5046c3..d56ca73adb4202 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/exploration/mv/PartitionCompensator.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/exploration/mv/PartitionCompensator.java @@ -218,11 +218,10 @@ private static Pair>, Pair>, Map>> invalidPartitions, - CascadesContext cascadesContext) { + public static boolean hasPartitionCompensation( + Pair>, Map>> invalidPartitions) { return invalidPartitions != null - && (!invalidPartitions.key().values().isEmpty() || !invalidPartitions.value().values().isEmpty()); + && (!invalidPartitions.key().isEmpty() || !invalidPartitions.value().isEmpty()); } /** diff --git a/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/exploration/mv/PartitionCompensatorTest.java b/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/exploration/mv/PartitionCompensatorTest.java index 1af0b983fee7a1..fa2977edabcc84 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/exploration/mv/PartitionCompensatorTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/exploration/mv/PartitionCompensatorTest.java @@ -41,6 +41,7 @@ import com.google.common.collect.ArrayListMultimap; import com.google.common.collect.ImmutableList; +import com.google.common.collect.ImmutableMap; import com.google.common.collect.ImmutableSet; import com.google.common.collect.Multimap; import org.junit.jupiter.api.Assertions; @@ -290,6 +291,20 @@ public void testNotNeedUnionRewriteWhenAllPartitions() throws Exception { Assertions.assertFalse(PartitionCompensator.needUnionRewrite(ctx, sc)); } + @Test + public void testHasPartitionCompensation() { + BaseTableInfo mvTableInfo = newBaseTableInfo(); + BaseColInfo baseColInfo = new BaseColInfo("c", newBaseTableInfo()); + + Assertions.assertFalse(PartitionCompensator.hasPartitionCompensation(null)); + Assertions.assertFalse(PartitionCompensator.hasPartitionCompensation( + Pair.of(ImmutableMap.of(), ImmutableMap.of()))); + Assertions.assertTrue(PartitionCompensator.hasPartitionCompensation( + Pair.of(ImmutableMap.of(mvTableInfo, ImmutableSet.of("mv_p1")), ImmutableMap.of()))); + Assertions.assertTrue(PartitionCompensator.hasPartitionCompensation( + Pair.of(ImmutableMap.of(), ImmutableMap.of(baseColInfo, ImmutableSet.of("base_p1"))))); + } + @Test public void testGetQueryUsedPartitionsAllAndPartial() { // Prepare qualifiers