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]

Reply via email to