This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new fa5667c12ab branch-4.1: [fix](mv) Compensate complete invalid roll-up
partition buckets #67882 (#68067)
fa5667c12ab is described below
commit fa5667c12ab0be4313580954d2794bdc745c0d85
Author: morrySnow <[email protected]>
AuthorDate: Wed Sep 16 23:30:14 2026 +0800
branch-4.1: [fix](mv) Compensate complete invalid roll-up partition buckets
#67882 (#68067)
### What problem does this PR solve?
Related PR: #67882
Problem Summary:
Backport the MV partition compensation fix to branch-4.1. When an
invalid roll-up MV partition intersects the query, compensate its
complete base-partition bucket. The unit test is adapted to the 4.1 API
signatures.
### Release note
None
### Check List (For Author)
- Test
- [x] Unit Test
- Behavior changed:
- [x] Yes. Invalid roll-up partition buckets are compensated completely.
- Does this need documentation?
- [x] No.
### Check List (For Reviewer who merge this PR)
- [ ] Confirm the release note
- [ ] Confirm test cases
- [ ] Confirm document
- [ ] Add branch pick label
---
.../rules/exploration/mv/PartitionCompensator.java | 7 +-
.../exploration/mv/PartitionCompensatorTest.java | 78 ++++++++++++++++++++++
2 files changed, 83 insertions(+), 2 deletions(-)
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 28a2eb6c55f..528babe6b04 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
@@ -198,8 +198,11 @@ public class PartitionCompensator {
// Base table partition maybe deleted, need not union
continue;
}
- Sets.intersection(baseTablePartitions,
queryUsedBaseTablePartitionNameSet)
- .copyInto(baseTableNeedUnionPartitionNameSet);
+ if (!Sets.intersection(baseTablePartitions,
queryUsedBaseTablePartitionNameSet).isEmpty()) {
+ // An MV partition is the atomic unit removed from the
rewritten plan. If any base
+ // partition in its roll-up bucket is used by the query,
compensate the whole bucket.
+ baseTableNeedUnionPartitionNameSet.addAll(baseTablePartitions);
+ }
}
// If related base table creates partitions or mv is created with ttl,
need base table union
Sets.difference(queryUsedBaseTablePartitionNameSet,
mvValidBaseTablePartitionNameSet)
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 1af0b983fee..60d99bfe4e9 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
@@ -486,6 +486,84 @@ public class PartitionCompensatorTest extends
TestWithFeService {
.forEach(v -> Assertions.assertEquals(expectedUnion, v));
}
+ @SuppressWarnings("unchecked")
+ @Test
+ public void
testCalcInvalidPartitionsCompensatesWholeIntersectingRollupBucket()
+ throws Exception {
+ DatabaseIf<?> baseDb = mockDatabase("cat", 1L, "db", 2L);
+ MTMVRelatedTableIf relatedTable = mockRelatedTableIf(
+ "base_t", 10L, ImmutableList.of("cat", "db", "base_t"),
baseDb);
+ BaseColInfo colInfo = new BaseColInfo("dt", new
BaseTableInfo(relatedTable));
+
+ DatabaseIf<?> mvDb = mockDatabase("internal", 3L, "mv_db", 4L);
+ MTMV mtmv = Mockito.mock(MTMV.class);
+ Mockito.when(mtmv.getName()).thenReturn("mv1");
+ Mockito.when(mtmv.getId()).thenReturn(100L);
+ Mockito.when(mtmv.getDatabase()).thenReturn(mvDb);
+ Mockito.when(mtmv.selectNonEmptyPartitionIds(ArgumentMatchers.any()))
+ .thenReturn(ImmutableList.of(1L));
+
+ long validMvPartitionId = 101L;
+ long partiallyStaleMvPartitionId = 102L;
+ long disjointMvPartitionId = 103L;
+ Partition validMvPartition = mockPartition(validMvPartitionId,
"mv_valid");
+ Partition partiallyStaleMvPartition =
mockPartition(partiallyStaleMvPartitionId, "mv_partially_stale");
+ Partition disjointMvPartition = mockPartition(disjointMvPartitionId,
"mv_disjoint");
+
Mockito.when(mtmv.getPartition(validMvPartitionId)).thenReturn(validMvPartition);
+
Mockito.when(mtmv.getPartition(partiallyStaleMvPartitionId)).thenReturn(partiallyStaleMvPartition);
+
Mockito.when(mtmv.getPartition(disjointMvPartitionId)).thenReturn(disjointMvPartition);
+
+ PartitionInfo mvPartitionInfo = Mockito.mock(PartitionInfo.class);
+ Mockito.when(mtmv.getPartitionInfo()).thenReturn(mvPartitionInfo);
+
Mockito.when(mvPartitionInfo.getType()).thenReturn(PartitionType.RANGE);
+ MTMVPartitionInfo mvPctInfo = Mockito.mock(MTMVPartitionInfo.class);
+ Mockito.when(mtmv.getMvPartitionInfo()).thenReturn(mvPctInfo);
+
Mockito.when(mvPctInfo.getPctTables()).thenReturn(ImmutableSet.of(relatedTable));
+
Mockito.when(mvPctInfo.getPctInfos()).thenReturn(ImmutableList.of(colInfo));
+
+ Map<String, Set<String>> relatedPartitionMapping = new HashMap<>();
+ relatedPartitionMapping.put("mv_valid", ImmutableSet.of("p3"));
+ relatedPartitionMapping.put("mv_partially_stale",
ImmutableSet.of("p1", "p2"));
+ relatedPartitionMapping.put("mv_disjoint", ImmutableSet.of("p4",
"p5"));
+ Map<MTMVRelatedTableIf, Map<String, Set<String>>> partitionMappings =
new HashMap<>();
+ partitionMappings.put(relatedTable, relatedPartitionMapping);
+
+ AsyncMaterializationContext matCtx =
Mockito.mock(AsyncMaterializationContext.class);
+ Mockito.when(matCtx.getMtmv()).thenReturn(mtmv);
+
Mockito.when(matCtx.calculatePartitionMappings()).thenReturn(partitionMappings);
+
+ Map<BaseTableInfo, Collection<Partition>> canRewriteMap = new
HashMap<>();
+ canRewriteMap.put(new BaseTableInfo(mtmv),
ImmutableList.of(validMvPartition));
+ StatementContext stmtCtx = Mockito.mock(StatementContext.class);
+
Mockito.when(stmtCtx.getMvCanRewritePartitionsMap()).thenReturn(canRewriteMap);
+ CascadesContext cascadesCtx = Mockito.mock(CascadesContext.class);
+ Mockito.when(cascadesCtx.getStatementContext()).thenReturn(stmtCtx);
+
+ LogicalOlapScan selectedMvScan = Mockito.mock(LogicalOlapScan.class);
+ Mockito.when(selectedMvScan.getTable()).thenReturn(mtmv);
+ Mockito.when(selectedMvScan.getSelectedPartitionIds())
+ .thenReturn(ImmutableList.of(
+ validMvPartitionId, partiallyStaleMvPartitionId,
disjointMvPartitionId));
+ Plan rewrittenPlan = Mockito.mock(Plan.class);
+ Mockito.when(rewrittenPlan.collectToList(ArgumentMatchers.any()))
+ .thenReturn(ImmutableList.of(selectedMvScan));
+
+ // The query touches p1 in the stale roll-up bucket and p3 in a valid
bucket. Once the stale
+ // MV partition is removed, its complete {p1, p2} mapping must be
compensated atomically.
+ // The completely disjoint {p4, p5} bucket is removed from the MV scan
but is not compensated.
+ Map<List<String>, Set<String>> queryUsedPartitions = new HashMap<>();
+ queryUsedPartitions.put(relatedTable.getFullQualifiers(),
ImmutableSet.of("p1", "p3"));
+
+ Pair<Map<BaseTableInfo, Set<String>>, Map<BaseColInfo, Set<String>>>
result =
+ PartitionCompensator.calcInvalidPartitions(
+ queryUsedPartitions, rewrittenPlan, matCtx,
cascadesCtx);
+
+ Assertions.assertNotNull(result);
+ Assertions.assertEquals(ImmutableSet.of("mv_partially_stale",
"mv_disjoint"),
+ result.key().get(new BaseTableInfo(mtmv)));
+ Assertions.assertEquals(ImmutableSet.of("p1", "p2"),
result.value().get(colInfo));
+ }
+
@SuppressWarnings("unchecked")
@Test
public void
testCalcInvalidPartitionsDoesNotCompensateBasePartitionsUnusedByQuery()
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]