github-actions[bot] commented on code in PR #68646:
URL: https://github.com/apache/doris/pull/68646#discussion_r4143552239


##########
fe/fe-core/src/main/java/org/apache/doris/alter/Alter.java:
##########
@@ -373,8 +374,30 @@ private boolean 
processAlterOlapTableInternal(List<AlterOp> alterOps, OlapTable
             throw new DdlException("Invalid alter operations: " + 
currentAlterOps);
         }
         if (needChangeMTMVState(alterOps)) {
-            Env.getCurrentEnv().getMtmvService()
-                .alterTable(oldBaseTableInfo, newBaseTableInfo, 
currentAlterOps.hasReplaceTableOp());
+            // Which columns an operation's effect on a view turns on is the 
operation's to say, see
+            // AlterOp#queryJudgedColumnNames, and every clause of the alter 
has to name them: a batch that
+            // mixes a dropped column with a type change is decided by neither 
-- no query says anything
+            // about a type change -- and stays invalidated the way it was 
before the queries were asked at
+            // all. Each of them also has to have reached the table. A schema 
change that is not a light one
+            // is applied by a job, which may not have run where this hook 
runs: the table still holds the
+            // column the change is about, every query still analyses against 
it, and an invalidation
+            // decided on that answer would be about the table from before the 
change. What is asked is
+            // whether the change has reached the table, which is the same 
fact the re-analysis reads, so
+            // the two answers cannot disagree.
+            boolean judgedByQuery = alterOps.stream().allMatch(op -> 
!op.queryJudgedColumnNames().isEmpty()
+                    && op.hasReachedTheTable(olapTable));

Review Comment:
   [P2] Read a stable column map before query judging an ALTER. Two concurrent 
light `ADD COLUMN` statements can interleave after A releases the table write 
lock: B's `rebuildFullSchema()` clears the plain `nameToColumn` TreeMap while 
A's `hasReachedTheTable()` reads it here (or in the later lambda). A sees its 
already-added column as absent, falls back to whole-MV invalidation, and drops 
the rewrite snapshot even when the MV reads neither new column. Hold the table 
read lock for each reached check or use a stable schema snapshot; keep analysis 
and journal waits outside the lock.



##########
fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelationManager.java:
##########
@@ -349,34 +359,167 @@ public void dropTable(Table table) {
         // because a dropped table is the one change whose query is gone 
beyond doubt. What the two record
         // is the same state either way. Unlike a rename it stays an 
invalidation: the table is gone for
         // good, so the state is not something a later alter can make obsolete.
-        processBaseTableChange(new BaseTableInfo(table), "The base table has 
been deleted:", false);
+        processBaseTableChange(new BaseTableInfo(table), "The base table has 
been deleted:", null);
     }
 
     /**
      * update mtmv status to `SCHEMA_CHANGE`.
      *
      * @param isReplace
+     * @param queryJudgedColumns the names the alter gives the table or takes 
away from it, which leave the
+     *                           judgement about each MV's state to that MV's 
own query, or null when the
+     *                           alter is not one a query decides. The names 
are carried rather than judged
+     *                           before the call because the judgement is 
about them; see
+     *                           {@code AlterOp#queryJudgedColumnNames} for 
which operations name one, and
+     *                           {@link #invalidateMvUnlessQueryHolds} for 
what is asked about it. A rename
+     *                           of the base table names no column: it is left 
to the record below, which
+     *                           says what the MV that keeps spelling the old 
name needs to hear
      */
     @Override
-    public void alterTable(BaseTableInfo oldTableInfo, Optional<BaseTableInfo> 
newTableInfo, boolean isReplace) {
+    public void alterTable(BaseTableInfo oldTableInfo, Optional<BaseTableInfo> 
newTableInfo, boolean isReplace,
+            QueryJudgedChange queryJudgedChange) {
         // when replace, need deal two table
         if (isReplace) {
             // REPLACE TABLE already invalidates the IVM baseline explicitly, 
see Alter#processReplaceTable
-            processBaseTableChange(newTableInfo.get(), "The base table has 
been updated:", false);
+            processBaseTableChange(newTableInfo.get(), "The base table has 
been updated:", null);
         }
-        boolean renamed = !isReplace && newTableInfo.isPresent()
-                && !Objects.equals(oldTableInfo.getTableName(), 
newTableInfo.get().getTableName());
-        // A rename is the one change whose query check is skipped: the MV 
query keeps spelling the old
-        // name, so it is unanalyzable by construction, and the reason it 
would be invalidated with --
-        // "the query is no longer analyzable" -- says less than the message 
this call records anyway.
-        boolean checkQueryUsable = !renamed;
-        processBaseTableChange(oldTableInfo, "The base table has been 
updated:", checkQueryUsable);
+        processBaseTableChange(oldTableInfo, "The base table has been 
updated:", queryJudgedChange);
     }
 
 
     /**
-     * An MV's query is only as good as the base table schema it was analyzed 
against. Re-analyzing the
-     * MV query here (right after the alter was applied) is what detects a 
changed column identity:
+     * Whether the query, as it is analysed now, reads a column of any of 
these names, and reads it where
+     * the change can reach it.
+     *
+     * <p>There are two places a name is the change's to answer for. One is a 
column of the table the change
+     * is about: that is the column this view's rows were computed from, and 
the names are matched
+     * case-insensitively because a name is what moves. The other is a column 
the query reaches across a
+     * scope boundary -- the plan records those on the Apply that stands for 
the subquery, whose correlation
+     * slots are the outer columns its right side reads -- because such a name 
is the scopes' to answer for
+     * rather than the query's: the nearest column to the reference answers 
for it, so a column the change
+     * takes away from a scope inside leaves the name to one outside, and a 
column it gives to a scope inside
+     * takes the name over. A name reached with the qualifier of another table 
inside the query's own scope
+     * is neither: no later change can move it, so one to a column it does not 
name is one this view's rows
+     * do not depend on.
+     */
+    private static boolean reachesAnyColumnOf(Plan plan, BaseTableInfo 
baseTableInfo, Set<String> columnNames) {
+        if (plan == null) {
+            // A query whose plan was not kept is one this cannot be answered 
about, and "it does" is the
+            // answer that keeps the view safe.
+            return true;
+        }
+        Set<String> names = Sets.newTreeSet(String.CASE_INSENSITIVE_ORDER);
+        names.addAll(columnNames);
+        LineageInfo lineage = LineageInfoExtractor.extractLineageInfo(plan);
+        for (SetMultimap<?, Expression> byType : 
lineage.getDirectLineageMap().values()) {
+            if (reachesAnyColumn(byType.values(), names, baseTableInfo)) {
+                return true;
+            }
+        }
+        // The dataset predicates once, not once per output column: the 
per-output copy of them the lineage
+        // also offers holds the same expressions for every column the query 
produces, and scanning it would
+        // visit each of them once per column.
+        if (reachesAnyColumn(lineage.getDatasetIndirectLineageMap().values(), 
names, baseTableInfo)) {
+            return true;
+        }
+        if (reachesAnyColumnOfASubquery(plan, names, baseTableInfo)) {
+            return true;
+        }
+        return reachesAnyColumnAcrossScopes(plan, lineage, names, 
baseTableInfo);
+    }
+
+    /** Whether this slot is a column of this table, through whatever views 
stand between the two. */
+    private static boolean isColumnOf(Slot slot, BaseTableInfo baseTableInfo) {
+        if (!(slot instanceof SlotReference)) {
+            return false;
+        }
+        return ((SlotReference) slot).getOriginalTable()
+                .map(table -> new BaseTableInfo(table).equals(baseTableInfo))
+                .orElse(false);
+    }
+
+    /**
+     * Whether a name the change is about is answered for inside a subquery, 
out of that subquery's own
+     * scope.
+     *
+     * <p>This is the one place a name can move without any column the view 
produces depending on it: the
+     * scope of a subquery is internal, so which column answers for a name 
there changes what the query
+     * returns -- a row, or none -- while every column of the view stays the 
one it was. The lineage of the
+     * view's columns does not reach it, so the scope the subquery became is 
read here, expression by
+     * expression, the way the lineage is read for the view's own.
+     *
+     * <p>Two things are read. One is a value the subquery itself names -- an 
expression of its own under one
+     * of these names, rather than a column of a table -- because that is what 
a name the change takes away
+     * falls back to, and it decides the rows whether the subquery is a 
predicate or a value. The other is a
+     * column of the table the change is about, which decides the rows only 
when the subquery's output is
+     * one the query reads: an EXISTS tests the rows of its subquery and not 
what it projects, so a name it
+     * projects and never compares is one this view's rows do not depend on.
+     */
+    private static boolean reachesAnyColumnOfASubquery(Plan plan, Set<String> 
names,
+            BaseTableInfo baseTableInfo) {
+        for (LogicalApply<?, ?> apply : 
plan.<LogicalApply>collectToList(LogicalApply.class::isInstance)) {
+            boolean outputDecidesRows = !((LogicalApply<?, ?>) 
apply).isExist();
+            for (Plan node : 
apply.right().<Plan>collectToList(Plan.class::isInstance)) {
+                for (Expression expression : node.getExpressions()) {
+                    if (isNamedByTheSubquery(expression, names)

Review Comment:
   [P2] Do not treat a coincidental subquery alias as a changed-table 
dependency. For a refreshed MV `SELECT o.id FROM outer_t o WHERE o.flag IN 
(SELECT 1 AS flag FROM inner_t i)`, a light `ADD COLUMN flag INT DEFAULT 0` on 
`inner_t` leaves the consumed literal `1` and all MV rows unchanged. The 
analyzed Apply right Project still contains `Alias(1, flag)`, so this name-only 
branch returns true and invalidates the MV, dropping its rewrite snapshot and 
forcing a whole refresh. Check whether a reference could actually fall back to 
this alias, rather than matching every alias named like the changed column.



##########
fe/fe-core/src/main/java/org/apache/doris/alter/Alter.java:
##########
@@ -373,8 +374,30 @@ private boolean 
processAlterOlapTableInternal(List<AlterOp> alterOps, OlapTable
             throw new DdlException("Invalid alter operations: " + 
currentAlterOps);
         }
         if (needChangeMTMVState(alterOps)) {
-            Env.getCurrentEnv().getMtmvService()
-                .alterTable(oldBaseTableInfo, newBaseTableInfo, 
currentAlterOps.hasReplaceTableOp());
+            // Which columns an operation's effect on a view turns on is the 
operation's to say, see
+            // AlterOp#queryJudgedColumnNames, and every clause of the alter 
has to name them: a batch that
+            // mixes a dropped column with a type change is decided by neither 
-- no query says anything
+            // about a type change -- and stays invalidated the way it was 
before the queries were asked at
+            // all. Each of them also has to have reached the table. A schema 
change that is not a light one
+            // is applied by a job, which may not have run where this hook 
runs: the table still holds the
+            // column the change is about, every query still analyses against 
it, and an invalidation
+            // decided on that answer would be about the table from before the 
change. What is asked is
+            // whether the change has reached the table, which is the same 
fact the re-analysis reads, so
+            // the two answers cannot disagree.
+            boolean judgedByQuery = alterOps.stream().allMatch(op -> 
!op.queryJudgedColumnNames().isEmpty()
+                    && op.hasReachedTheTable(olapTable));
+            // The names of those columns go to the hook rather than a 
verdict: what the judgement is about
+            // is the column, and the hook holds the query's answer against it 
-- both while asking, in case
+            // the query can reach the name some other way now, and once it 
has answered, in case the table
+            // is no longer the one that answered. See MTMVRelationManager.
+            MTMVHookService.QueryJudgedChange queryJudgedChange = judgedByQuery
+                    ? new MTMVHookService.QueryJudgedChange(
+                            
alterOps.stream().map(AlterOp::queryJudgedColumnNames).flatMap(Set::stream)
+                                    .collect(Collectors.toSet()),
+                            () -> alterOps.stream().allMatch(op -> 
op.hasReachedTheTable(olapTable)))
+                    : null;
+            Env.getCurrentEnv().getMtmvService().alterTable(oldBaseTableInfo, 
newBaseTableInfo,

Review Comment:
   [P1] Close the gap between publishing a light ADD and invalidating dependent 
MVs. `SchemaChangeHandler.process` releases the base-table write lock after 
installing the new column, and this call only then re-analyzes each MV before 
invalidating it. A concurrent query can bind the new `inner_t.flag` in `o.flag 
IN (SELECT flag FROM inner_t)` while the refreshed MV still holds rows from the 
old outer binding and remains NORMAL with unchanged base visible versions. With 
supported `mtmv_cache_manage_num=0` (or a cache miss), rewrite builds the MV 
plan against that new binding and can return a stale row. Establish an 
invalidation barrier before the new schema is visible to query planners.



##########
fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelationManager.java:
##########
@@ -349,34 +359,167 @@ public void dropTable(Table table) {
         // because a dropped table is the one change whose query is gone 
beyond doubt. What the two record
         // is the same state either way. Unlike a rename it stays an 
invalidation: the table is gone for
         // good, so the state is not something a later alter can make obsolete.
-        processBaseTableChange(new BaseTableInfo(table), "The base table has 
been deleted:", false);
+        processBaseTableChange(new BaseTableInfo(table), "The base table has 
been deleted:", null);
     }
 
     /**
      * update mtmv status to `SCHEMA_CHANGE`.
      *
      * @param isReplace
+     * @param queryJudgedColumns the names the alter gives the table or takes 
away from it, which leave the
+     *                           judgement about each MV's state to that MV's 
own query, or null when the
+     *                           alter is not one a query decides. The names 
are carried rather than judged
+     *                           before the call because the judgement is 
about them; see
+     *                           {@code AlterOp#queryJudgedColumnNames} for 
which operations name one, and
+     *                           {@link #invalidateMvUnlessQueryHolds} for 
what is asked about it. A rename
+     *                           of the base table names no column: it is left 
to the record below, which
+     *                           says what the MV that keeps spelling the old 
name needs to hear
      */
     @Override
-    public void alterTable(BaseTableInfo oldTableInfo, Optional<BaseTableInfo> 
newTableInfo, boolean isReplace) {
+    public void alterTable(BaseTableInfo oldTableInfo, Optional<BaseTableInfo> 
newTableInfo, boolean isReplace,
+            QueryJudgedChange queryJudgedChange) {
         // when replace, need deal two table
         if (isReplace) {
             // REPLACE TABLE already invalidates the IVM baseline explicitly, 
see Alter#processReplaceTable
-            processBaseTableChange(newTableInfo.get(), "The base table has 
been updated:", false);
+            processBaseTableChange(newTableInfo.get(), "The base table has 
been updated:", null);
         }
-        boolean renamed = !isReplace && newTableInfo.isPresent()
-                && !Objects.equals(oldTableInfo.getTableName(), 
newTableInfo.get().getTableName());
-        // A rename is the one change whose query check is skipped: the MV 
query keeps spelling the old
-        // name, so it is unanalyzable by construction, and the reason it 
would be invalidated with --
-        // "the query is no longer analyzable" -- says less than the message 
this call records anyway.
-        boolean checkQueryUsable = !renamed;
-        processBaseTableChange(oldTableInfo, "The base table has been 
updated:", checkQueryUsable);
+        processBaseTableChange(oldTableInfo, "The base table has been 
updated:", queryJudgedChange);
     }
 
 
     /**
-     * An MV's query is only as good as the base table schema it was analyzed 
against. Re-analyzing the
-     * MV query here (right after the alter was applied) is what detects a 
changed column identity:
+     * Whether the query, as it is analysed now, reads a column of any of 
these names, and reads it where
+     * the change can reach it.
+     *
+     * <p>There are two places a name is the change's to answer for. One is a 
column of the table the change
+     * is about: that is the column this view's rows were computed from, and 
the names are matched
+     * case-insensitively because a name is what moves. The other is a column 
the query reaches across a
+     * scope boundary -- the plan records those on the Apply that stands for 
the subquery, whose correlation
+     * slots are the outer columns its right side reads -- because such a name 
is the scopes' to answer for
+     * rather than the query's: the nearest column to the reference answers 
for it, so a column the change
+     * takes away from a scope inside leaves the name to one outside, and a 
column it gives to a scope inside
+     * takes the name over. A name reached with the qualifier of another table 
inside the query's own scope
+     * is neither: no later change can move it, so one to a column it does not 
name is one this view's rows
+     * do not depend on.
+     */
+    private static boolean reachesAnyColumnOf(Plan plan, BaseTableInfo 
baseTableInfo, Set<String> columnNames) {
+        if (plan == null) {
+            // A query whose plan was not kept is one this cannot be answered 
about, and "it does" is the
+            // answer that keeps the view safe.
+            return true;
+        }
+        Set<String> names = Sets.newTreeSet(String.CASE_INSENSITIVE_ORDER);
+        names.addAll(columnNames);
+        LineageInfo lineage = LineageInfoExtractor.extractLineageInfo(plan);
+        for (SetMultimap<?, Expression> byType : 
lineage.getDirectLineageMap().values()) {
+            if (reachesAnyColumn(byType.values(), names, baseTableInfo)) {
+                return true;
+            }
+        }
+        // The dataset predicates once, not once per output column: the 
per-output copy of them the lineage
+        // also offers holds the same expressions for every column the query 
produces, and scanning it would
+        // visit each of them once per column.
+        if (reachesAnyColumn(lineage.getDatasetIndirectLineageMap().values(), 
names, baseTableInfo)) {

Review Comment:
   [P2] Ignore lineage from CTE producers the MV result never consumes. When 
`inner_t` initially lacks `flag`, a refreshed MV `WITH unused AS (SELECT 1 AS 
flag, COUNT(*) AS n FROM inner_t i GROUP BY flag HAVING flag=1) SELECT o.id 
FROM outer_t o` has no dependency of its result on that producer. A light `ADD 
COLUMN inner_t.flag INT DEFAULT 0` rebinds the unused producer's GROUP BY 
`flag` to `i.flag`, while every MV row stays the same. The analyzed plan 
retains that producer under a CTE anchor, so its `i.flag` enters dataset 
lineage here and triggers whole-MV invalidation and snapshot loss. Restrict 
this scan to producers reachable from the result.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to