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 2871c53ed1e branch-4.1:[fix](mtmv) Refresh MTMV after excluded trigger
tables change (#64041) (#67267)
2871c53ed1e is described below
commit 2871c53ed1eb5c5bca53870c20d71f4d2fd4fd1a
Author: seawinde <[email protected]>
AuthorDate: Fri Aug 28 20:56:13 2026 +0800
branch-4.1:[fix](mtmv) Refresh MTMV after excluded trigger tables change
(#64041) (#67267)
pr: #64041
commitId: 3d5cee2481e
---
.../main/java/org/apache/doris/catalog/MTMV.java | 47 +++++++--
.../apache/doris/job/extensions/mtmv/MTMVTask.java | 5 +
.../java/org/apache/doris/mtmv/MTMVTaskTest.java | 52 ++++++++++
.../test/java/org/apache/doris/mtmv/MTMVTest.java | 105 +++++++++++++++++++++
.../data/mtmv_p0/test_create_mtmv_with_view.out | 3 +-
.../mtmv_p0/test_excluded_trigger_table_mtmv.out | 6 ++
.../test_excluded_trigger_table_mtmv.groovy | 6 +-
7 files changed, 211 insertions(+), 13 deletions(-)
diff --git a/fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java
b/fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java
index dc03e15dd0f..82955035282 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java
@@ -273,7 +273,21 @@ public class MTMV extends OlapTable {
public Map<String, String> alterMvProperties(Map<String, String>
mvProperties) {
writeMvLock();
try {
+ boolean containsExcludedTriggerTables = mvProperties.containsKey(
+ PropertyAnalyzer.PROPERTIES_EXCLUDED_TRIGGER_TABLES);
+ Set<TableName> oldExcludedTriggerTables =
containsExcludedTriggerTables
+ ? parseExcludedTriggerTables()
+ : Sets.newHashSet();
this.mvProperties.putAll(mvProperties);
+ if (containsExcludedTriggerTables) {
+ Set<TableName> newExcludedTriggerTables =
parseExcludedTriggerTables();
+ if
(!oldExcludedTriggerTables.equals(newExcludedTriggerTables)) {
+ // excluded_trigger_tables changes the refresh baseline
semantics. Invalidate the old
+ // snapshots so the next AUTO refresh rebuilds a complete
baseline with the new rules.
+ this.schemaChangeVersion++;
+ this.refreshSnapshot = new MTMVRefreshSnapshot();
+ }
+ }
return this.mvProperties;
} finally {
writeMvUnlock();
@@ -334,22 +348,26 @@ public class MTMV extends OlapTable {
}
public Set<TableName> getExcludedTriggerTables() {
- Set<TableName> res = Sets.newHashSet();
readMvLock();
try {
- if
(StringUtils.isEmpty(mvProperties.get(PropertyAnalyzer.PROPERTIES_EXCLUDED_TRIGGER_TABLES)))
{
- return res;
- }
- String[] split =
mvProperties.get(PropertyAnalyzer.PROPERTIES_EXCLUDED_TRIGGER_TABLES).split(",");
- for (String alias : split) {
- res.add(new TableName(alias));
- }
- return res;
+ return parseExcludedTriggerTables();
} finally {
readMvUnlock();
}
}
+ private Set<TableName> parseExcludedTriggerTables() {
+ Set<TableName> res = Sets.newHashSet();
+ if
(StringUtils.isEmpty(mvProperties.get(PropertyAnalyzer.PROPERTIES_EXCLUDED_TRIGGER_TABLES)))
{
+ return res;
+ }
+ String[] split =
mvProperties.get(PropertyAnalyzer.PROPERTIES_EXCLUDED_TRIGGER_TABLES).split(",");
+ for (String alias : split) {
+ res.add(new TableName(alias));
+ }
+ return res;
+ }
+
public Set<TableName> getQueryRewriteConsistencyRelaxedTables() {
Set<TableName> res = Sets.newHashSet();
readMvLock();
@@ -437,6 +455,17 @@ public class MTMV extends OlapTable {
return refreshSnapshot;
}
+ public boolean hasCompleteRefreshSnapshot() {
+ Set<String> partitionNames = getPartitionNames();
+ readMvLock();
+ try {
+ // A refresh baseline is complete only when every current MV
partition has a snapshot.
+ return
refreshSnapshot.getPartitionSnapshots().keySet().containsAll(partitionNames);
+ } finally {
+ readMvUnlock();
+ }
+ }
+
public long getSchemaChangeVersion() {
readMvLock();
try {
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java
index e71798f21b5..f3dbabf97c9 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java
@@ -688,6 +688,11 @@ public class MTMVTask extends AbstractTask {
if (mtmv.getRefreshInfo().getRefreshMethod() ==
RefreshMethod.COMPLETE) {
return Lists.newArrayList(mtmv.getPartitionNames());
}
+ // An incomplete baseline cannot be checked by isMTMVSync, because the
current exclude rules may
+ // skip the changed base tables and incorrectly mark the MV as fresh.
Rebuild it with a full refresh.
+ if (!mtmv.hasCompleteRefreshSnapshot()) {
+ return Lists.newArrayList(mtmv.getPartitionNames());
+ }
// check if data is fresh
// We need to use a newly generated relationship and cannot retrieve
it using mtmv.getRelation()
// to avoid rebuilding the baseTable and causing a change in the
tableId
diff --git a/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVTaskTest.java
b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVTaskTest.java
index 44b88cc506f..8d96489bf93 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVTaskTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVTaskTest.java
@@ -39,6 +39,7 @@ import com.google.common.collect.Lists;
import com.google.common.collect.Sets;
import mockit.Expectations;
import mockit.Mocked;
+import mockit.Verifications;
import org.apache.commons.collections4.CollectionUtils;
import org.junit.Assert;
import org.junit.Before;
@@ -103,6 +104,10 @@ public class MTMVTaskTest {
mtmvRefreshInfo.getRefreshMethod();
minTimes = 0;
result = RefreshMethod.COMPLETE;
+
+ mtmv.hasCompleteRefreshSnapshot();
+ minTimes = 0;
+ result = true;
}
};
}
@@ -147,6 +152,53 @@ public class MTMVTaskTest {
Assert.assertEquals(allPartitionNames, result);
}
+ @Test
+ public void
testCalculateNeedRefreshPartitionsSystemIncompleteRefreshSnapshot() throws
AnalysisException {
+ new Expectations() {
+ {
+ mtmvRefreshInfo.getRefreshMethod();
+ minTimes = 0;
+ result = RefreshMethod.AUTO;
+
+ mtmv.hasCompleteRefreshSnapshot();
+ minTimes = 0;
+ result = false;
+ }
+ };
+
+ MTMVTaskContext context = new
MTMVTaskContext(MTMVTaskTriggerMode.SYSTEM);
+ MTMVTask task = new MTMVTask(mtmv, relation, context);
+ List<String> result = task.calculateNeedRefreshPartitions(null);
+
+ Assert.assertEquals(allPartitionNames, result);
+ new Verifications() {
+ {
+ mtmvPartitionUtil.isMTMVSync((MTMVRefreshContext) any,
(Set<BaseTableInfo>) any,
+ (Set<TableName>) any);
+ times = 0;
+ }
+ };
+ }
+
+ @Test
+ public void
testCalculateNeedRefreshPartitionsManualPartitionsIncompleteRefreshSnapshot()
+ throws AnalysisException {
+ new Expectations() {
+ {
+ mtmv.hasCompleteRefreshSnapshot();
+ minTimes = 0;
+ result = false;
+ }
+ };
+
+ MTMVTaskContext context = new
MTMVTaskContext(MTMVTaskTriggerMode.MANUAL, Lists.newArrayList(poneName),
+ false, null);
+ MTMVTask task = new MTMVTask(mtmv, relation, context);
+ List<String> result = task.calculateNeedRefreshPartitions(null);
+
+ Assert.assertEquals(Lists.newArrayList(poneName), result);
+ }
+
@Test
public void testCalculateNeedRefreshPartitionsSystemNotSyncComplete()
throws AnalysisException {
new Expectations() {
diff --git a/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVTest.java
b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVTest.java
index 0f9f2e4a0ed..b25e84c3ffe 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVTest.java
@@ -177,6 +177,111 @@ public class MTMVTest {
Assert.assertTrue(excludedTriggerTables.contains(new TableName(null,
null, "t3")));
}
+ @Test
+ public void testAlterMvPropertiesWithExcludedTriggerTablesChange() {
+ Map<String, String> mvProperties = Maps.newHashMap();
+ mvProperties.put(PropertyAnalyzer.PROPERTIES_EXCLUDED_TRIGGER_TABLES,
"t1");
+ MTMV mtmv = new MTMV();
+ mtmv.setMvProperties(mvProperties);
+ MTMVStatus status = new MTMVStatus(MTMVState.NORMAL, null);
+ mtmv.setStatus(status);
+ MTMVRefreshSnapshot refreshSnapshot = new MTMVRefreshSnapshot();
+ refreshSnapshot.getPartitionSnapshots().put("p1", new
MTMVRefreshPartitionSnapshot());
+ mtmv.setRefreshSnapshot(refreshSnapshot);
+
+ long oldSchemaChangeVersion = mtmv.getSchemaChangeVersion();
+ Map<String, String> newProperties = Maps.newHashMap();
+ newProperties.put(PropertyAnalyzer.PROPERTIES_EXCLUDED_TRIGGER_TABLES,
"db1.t1");
+
+ mtmv.alterMvProperties(newProperties);
+
+ Assert.assertEquals(MTMVState.NORMAL, mtmv.getStatus().getState());
+ Assert.assertEquals(oldSchemaChangeVersion + 1,
mtmv.getSchemaChangeVersion());
+
Assert.assertTrue(mtmv.getRefreshSnapshot().getPartitionSnapshots().isEmpty());
+
+ mtmv.getRefreshSnapshot().getPartitionSnapshots().put("p1", new
MTMVRefreshPartitionSnapshot());
+ oldSchemaChangeVersion = mtmv.getSchemaChangeVersion();
+ newProperties.put(PropertyAnalyzer.PROPERTIES_EXCLUDED_TRIGGER_TABLES,
"internal.db1.t1");
+
+ mtmv.alterMvProperties(newProperties);
+
+ Assert.assertEquals(MTMVState.NORMAL, mtmv.getStatus().getState());
+ Assert.assertEquals(oldSchemaChangeVersion + 1,
mtmv.getSchemaChangeVersion());
+
Assert.assertTrue(mtmv.getRefreshSnapshot().getPartitionSnapshots().isEmpty());
+ }
+
+ @Test
+ public void testAlterMvPropertiesWithSameExcludedTriggerTables() {
+ Map<String, String> mvProperties = Maps.newHashMap();
+ mvProperties.put(PropertyAnalyzer.PROPERTIES_EXCLUDED_TRIGGER_TABLES,
"t1,t2");
+ MTMV mtmv = new MTMV();
+ mtmv.setMvProperties(mvProperties);
+ MTMVRefreshSnapshot refreshSnapshot = new MTMVRefreshSnapshot();
+ refreshSnapshot.getPartitionSnapshots().put("p1", new
MTMVRefreshPartitionSnapshot());
+ mtmv.setRefreshSnapshot(refreshSnapshot);
+
+ long oldSchemaChangeVersion = mtmv.getSchemaChangeVersion();
+ Map<String, String> newProperties = Maps.newHashMap();
+ newProperties.put(PropertyAnalyzer.PROPERTIES_EXCLUDED_TRIGGER_TABLES,
"t2,t1");
+
+ mtmv.alterMvProperties(newProperties);
+
+ Assert.assertEquals(oldSchemaChangeVersion,
mtmv.getSchemaChangeVersion());
+
Assert.assertFalse(mtmv.getRefreshSnapshot().getPartitionSnapshots().isEmpty());
+ }
+
+ @Test
+ public void testAlterMvPropertiesWithReducedExcludedTriggerTables() {
+ Map<String, String> mvProperties = Maps.newHashMap();
+ mvProperties.put(PropertyAnalyzer.PROPERTIES_EXCLUDED_TRIGGER_TABLES,
"t1,t2");
+ MTMV mtmv = new MTMV();
+ mtmv.setMvProperties(mvProperties);
+ mtmv.setStatus(new MTMVStatus(MTMVState.NORMAL, null));
+ MTMVRefreshSnapshot refreshSnapshot = new MTMVRefreshSnapshot();
+ refreshSnapshot.getPartitionSnapshots().put("p1", new
MTMVRefreshPartitionSnapshot());
+ mtmv.setRefreshSnapshot(refreshSnapshot);
+
+ long oldSchemaChangeVersion = mtmv.getSchemaChangeVersion();
+ Map<String, String> newProperties = Maps.newHashMap();
+ newProperties.put(PropertyAnalyzer.PROPERTIES_EXCLUDED_TRIGGER_TABLES,
"t1");
+
+ mtmv.alterMvProperties(newProperties);
+
+ Assert.assertEquals(MTMVState.NORMAL, mtmv.getStatus().getState());
+ Assert.assertEquals(oldSchemaChangeVersion + 1,
mtmv.getSchemaChangeVersion());
+
Assert.assertTrue(mtmv.getRefreshSnapshot().getPartitionSnapshots().isEmpty());
+
+ mtmv.getRefreshSnapshot().getPartitionSnapshots().put("p1", new
MTMVRefreshPartitionSnapshot());
+ oldSchemaChangeVersion = mtmv.getSchemaChangeVersion();
+ newProperties.put(PropertyAnalyzer.PROPERTIES_EXCLUDED_TRIGGER_TABLES,
"");
+
+ mtmv.alterMvProperties(newProperties);
+
+ Assert.assertEquals(MTMVState.NORMAL, mtmv.getStatus().getState());
+ Assert.assertEquals(oldSchemaChangeVersion + 1,
mtmv.getSchemaChangeVersion());
+
Assert.assertTrue(mtmv.getRefreshSnapshot().getPartitionSnapshots().isEmpty());
+ }
+
+ @Test
+ public void testAlterMvPropertiesWithOtherProperty() {
+ Map<String, String> mvProperties = Maps.newHashMap();
+ mvProperties.put(PropertyAnalyzer.PROPERTIES_EXCLUDED_TRIGGER_TABLES,
"t1");
+ MTMV mtmv = new MTMV();
+ mtmv.setMvProperties(mvProperties);
+ MTMVRefreshSnapshot refreshSnapshot = new MTMVRefreshSnapshot();
+ refreshSnapshot.getPartitionSnapshots().put("p1", new
MTMVRefreshPartitionSnapshot());
+ mtmv.setRefreshSnapshot(refreshSnapshot);
+
+ long oldSchemaChangeVersion = mtmv.getSchemaChangeVersion();
+ Map<String, String> newProperties = Maps.newHashMap();
+ newProperties.put(PropertyAnalyzer.PROPERTIES_GRACE_PERIOD, "10");
+
+ mtmv.alterMvProperties(newProperties);
+
+ Assert.assertEquals(oldSchemaChangeVersion,
mtmv.getSchemaChangeVersion());
+
Assert.assertFalse(mtmv.getRefreshSnapshot().getPartitionSnapshots().isEmpty());
+ }
+
@Test
public void testAlterStatus() {
MTMV mtmv = new MTMV();
diff --git a/regression-test/data/mtmv_p0/test_create_mtmv_with_view.out
b/regression-test/data/mtmv_p0/test_create_mtmv_with_view.out
index 90c7a97c096..f083632719a 100644
--- a/regression-test/data/mtmv_p0/test_create_mtmv_with_view.out
+++ b/regression-test/data/mtmv_p0/test_create_mtmv_with_view.out
@@ -22,9 +22,10 @@ COMPLETE
2 2
-- !trigger_table_not_need_refresh --
-NOT_REFRESH
+COMPLETE
-- !after_trigger_table --
1 1
2 2
+3 3
diff --git a/regression-test/data/mtmv_p0/test_excluded_trigger_table_mtmv.out
b/regression-test/data/mtmv_p0/test_excluded_trigger_table_mtmv.out
index 4a2ede2fe66..4c0bc1f0076 100644
--- a/regression-test/data/mtmv_p0/test_excluded_trigger_table_mtmv.out
+++ b/regression-test/data/mtmv_p0/test_excluded_trigger_table_mtmv.out
@@ -4,12 +4,18 @@
-- !true_table --
1 1
+2 2
-- !true_db_table --
1 1
+2 2
+3 3
-- !true_ctl_db_table --
1 1
+2 2
+3 3
+4 4
-- !false_ctl_db_table --
1 1
diff --git
a/regression-test/suites/mtmv_p0/test_excluded_trigger_table_mtmv.groovy
b/regression-test/suites/mtmv_p0/test_excluded_trigger_table_mtmv.groovy
index 6a2264be699..4794470cd3c 100644
--- a/regression-test/suites/mtmv_p0/test_excluded_trigger_table_mtmv.groovy
+++ b/regression-test/suites/mtmv_p0/test_excluded_trigger_table_mtmv.groovy
@@ -65,7 +65,7 @@ suite("test_excluded_trigger_table_mtmv","mtmv") {
REFRESH MATERIALIZED VIEW ${mvName} AUTO
"""
waitingMTMVTaskFinishedByMvName(mvName)
- // should not refresh
+ // should refresh because excluded_trigger_tables changed and refresh
baseline should be rebuilt
order_qt_true_table "SELECT * FROM ${mvName}"
sql """
@@ -78,7 +78,7 @@ suite("test_excluded_trigger_table_mtmv","mtmv") {
REFRESH MATERIALIZED VIEW ${mvName} AUTO
"""
waitingMTMVTaskFinishedByMvName(mvName)
- // should not refresh
+ // should refresh because excluded_trigger_tables changed and refresh
baseline should be rebuilt
order_qt_true_db_table "SELECT * FROM ${mvName}"
sql """
@@ -91,7 +91,7 @@ suite("test_excluded_trigger_table_mtmv","mtmv") {
REFRESH MATERIALIZED VIEW ${mvName} AUTO
"""
waitingMTMVTaskFinishedByMvName(mvName)
- // should not refresh
+ // should refresh because excluded_trigger_tables changed and refresh
baseline should be rebuilt
order_qt_true_ctl_db_table "SELECT * FROM ${mvName}"
sql """
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]