yujun777 commented on code in PR #68390:
URL: https://github.com/apache/doris/pull/68390#discussion_r4122827607


##########
fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java:
##########
@@ -715,13 +796,52 @@ private AttemptResultType 
executeIvmAttempt(MTMVRefreshContext refreshContext,
                     + "Continuing with COMPLETE refresh.", mtmv.getName(), 
getTaskId());
             return AttemptResultType.FALLBACK_TO_COMPLETE;
         }
+        // The partitions the criterion says must be rebuilt rather than 
caught up: the delta path can only
+        // append, so a partition it treated as current would record that in 
its epoch while its rows still
+        // come from before the change. Rebuilt first, with the partition 
executor, because that is the
+        // full recomputation they need -- and only in this task's batches, so 
a change that arrives while
+        // it runs leaves them dirty for the next round instead of being 
swallowed.
+        // One read of the states decides both what has to be rebuilt and the 
requirement each batch may
+        // write back. Reading them separately would leave a window between 
the two in which a mark lands,
+        // the routing decision does not see it, and the batch that follows 
captures the raised requirement
+        // and records it as met by a delta that cannot remove the rows that 
mark made unusable.
+        Map<String, MTMVPartitionState> plannedStates = 
mtmv.getPartitionStates();
+        Set<String> livePartitionNames = mtmv.getPartitionNames();
+        Set<String> dirtyPartitions = Sets.newLinkedHashSet();
+        Map<String, Long> plannedEpochs = Maps.newHashMap();
+        for (Entry<String, MTMVPartitionState> plannedState : 
plannedStates.entrySet()) {
+            if (!livePartitionNames.contains(plannedState.getKey())) {
+                continue;
+            }
+            plannedEpochs.put(plannedState.getKey(), 
plannedState.getValue().getLatestEpoch());
+            if (plannedState.getValue().isDirty()) {
+                dirtyPartitions.add(plannedState.getKey());
+            }
+        }
+        this.ivmPlannedEpochs = plannedEpochs;
+        if (!dirtyPartitions.isEmpty()) {
+            LOG.info("Rebuilding {} invalidated MV partitions before the 
incremental refresh, mv={}, taskId={}",
+                    dirtyPartitions.size(), mtmv.getName(), getTaskId());
+            List<String> toRebuild = Lists.newArrayList(dirtyPartitions);
+            toRebuild.sort(Comparator.naturalOrder());
+            this.refreshMode = generateRefreshMode(toRebuild);
+            try {
+                executePartitionBasedRefresh(refreshContext, 
RefreshMode.PARTITIONS, ctx, toRebuild);
+            } finally {
+                recordRebuiltPartitions(request);
+            }
+        }
         MTMVRefreshContext currentRefreshContext = refreshContext;
         int ivmAttemptLimit = Math.max(Config.max_query_retry_time, 0) + 1;
         IvmIncrRefreshResult ivmResult = null;
         for (int partitionSyncRetryCount = 0;
                 partitionSyncRetryCount < ivmAttemptLimit; 
partitionSyncRetryCount++) {
-            ivmResult = executeSingleIvmAttempt(currentRefreshContext);
+            ivmResult = executeSingleIvmAttempt(currentRefreshContext, 
dirtyPartitions);
             if (ivmResult.isSuccess()) {
+                // The delta has run, and it is what brought the tables those 
partitions were rebuilt from up
+                // to date: its target is the partitions the change touches, 
which includes the ones the
+                // rebuild replaced. What the rebuild wrote is current now, so 
what was held can be recorded.
+                redeemHeldRecords();

Review Comment:
   Right, and the step it skips is exactly what the publication is for. 
`redeemHeldRecords` runs on `ivmResult.isSuccess()`, and 
`executeSingleIvmAttempt` returns success without reaching `doRefresh` when its 
scope comes out empty -- and the scope is *what the attempt records*, emptied 
by taking the partitions the rebuild replaced out of it, not a statement that 
there is nothing left to consume. It does not even need the retry to happen: 
the retry is simply the shape the code's own test already pins as reachable, 
since `adoptPartitionsCreatedByTheRetry` adds the recreated partition to the 
dirty set, so the scope can empty on the second attempt while the first 
attempt's delta failed.
   
   `42ed8c44575` runs the delta for the records it holds back: the skip now 
needs the scope to be empty *and* nothing held. `hasRecordsHeldForTheDelta` 
states why the two questions differ, and why the delta does not need the scope 
to consume the change -- its plan is the MV's own query over the streams, and 
the partitions it writes are the ones the change reaches, which includes the 
partitions the rebuild replaced. The attempt still records nothing (the scope 
is what it records), and what was held is published by the delta's success.
   
   The unit test pins both directions: 
`MTMVTaskTest#testTheIncrementalAttemptRunsForTheRecordsAPartialReadHeldBack` 
drives the state you describe -- held records, every refresh-needing partition 
dirty, so the scope is empty -- and asserts that `doRefresh` ran and that the 
held records were then published; `...IsSkippedWhenThereIsNothingHeldBack` pins 
that an attempt with nothing held is still skipped.
   
   The end-to-end retry-to-empty sequence is still not built: it needs a 
`MV_PARTITION_NOT_FOUND` that the following sync resolves, and I have not found 
a recipe that induces it reliably rather than by timing. Say the word and I 
will add it as its own suite rather than pinning the mechanism at the unit 
level only.
   



##########
fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java:
##########
@@ -444,7 +499,33 @@ private RefreshRequest resolveRefreshRequest() throws 
JobException {
                 Lists.newArrayList(), false);
     }
 
-    private List<RefreshAttemptType> buildAttempts(RefreshRequest request, 
boolean containsOneRowRelation) {
+    private List<RefreshAttemptType> buildAttempts(RefreshRequest request, 
boolean containsOneRowRelation)
+            throws JobException {
+        // A schema-level invalidation is not a set of dirty partitions: it 
means every partition, including
+        // the ones partition sync has not created yet, and no per-partition 
requirement can express that.
+        // IVM only -- a non-IVM MV reaches the same effect through its 
cleared snapshot, which its own
+        // refresh already depends on.
+        //
+        // Judged before the initial-refresh shortcut below, which also 
answers COMPLETE: an MV that has
+        // never been refreshed and reads an excluded trigger table (or a 
one-row relation) has to be built
+        // by a whole-MV refresh, and a request that may not fall back has to 
hear that rather than have it
+        // decided for it -- otherwise the refusal here would be unreachable 
in exactly the state it names.
+        if (mtmv.isIvm() && !request.explicitPartitions
+                && mtmv.getStatus().getState() == MTMVState.SCHEMA_CHANGE) {
+            if (request.refreshMode == RefreshMode.PARTITIONS && 
!request.allowFallback) {

Review Comment:
   This one is by design, and the widening is reported rather than silent: the 
task's `RefreshMode` reads COMPLETE after a strict `REFRESH ... INCREMENTAL` 
met an invalidated MV. SCHEMA_CHANGE is a property of the MV, not a scope the 
request named -- its own baseline is gone, and its query may no longer analyze 
at all (a rename is what the suite uses) -- so nothing incremental can repair 
it, and the COMPLETE refresh restores exactly the baseline the incremental 
attempt would have read. Refusing would leave a strict caller in a state that 
only a COMPLETE can leave, which is what the message on the PARTITIONS branch 
already tells them to run.
   
   That is a different question from the one the all-dirty and unusable-stream 
shortcuts answer. Those guard a strict request against a widening it did not 
ask for -- rebuilding partitions it excluded, or a stream reset its incremental 
attempt would have failed on instead. Here the state, not the request's scope, 
is what forces the rebuild, and the result says so.
   
   
`MTMVTaskTest#testBuildAttemptsRebuildsTheWholeMvForAStrictIncrementalInSchemaChange`
 and case 3 of `test_ivm_partition_epoch_rebuild` pin it, including the rebuilt 
count a request that did not ask for the rebuild reports.
   



##########
fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java:
##########
@@ -738,28 +1067,61 @@ public boolean invalidateIvmBaseline(BaseTableInfo 
baseTableInfo, Map<String, Lo
                     + "changedPartitions={}", name, baseTableInfo, 
changedPartitions);
             return false;
         }
+        if (!affectedMvPartitions.isPresent()) {
+            // A narrower rebuild could leave a partition holding rows of the 
changed base partition
+            // untouched, and those rows cannot be repaired later: the change 
emitted no row binlog. The
+            // whole MV is invalidated instead, which says "every partition, 
including the ones partition
+            // sync has not created yet" -- what a per-partition requirement 
cannot express.
+            invalidateWholeMv(reason).await();
+            return true;
+        }
         EditLogItem editLogItem;
         writeMvLock();
         try {
-            if (ivmInfo == null) {
-                ivmInfo = new IvmInfo();
-            }
-            if (!affectedMvPartitions.isPresent()) {
-                // A narrower rebuild could leave a partition holding rows of 
the changed base partition
-                // untouched, and those rows cannot be repaired later: the 
change emitted no row binlog.
-                ivmInfo.requireCompleteBaselineRebuild();
-            } else {
-                
ivmInfo.addPendingBaselineRebuildPartitions(affectedMvPartitions.get());
+            // Placed: the partitions that read the change get the requirement 
raised, which is what sends
+            // them to a rebuild while every other partition keeps catching up 
incrementally. No version
+            // bump here -- a partial invalidation does not invalidate a task 
result, and the requirement it
+            // raises survives the write-back by construction.
+            Set<String> marked = 
markIvmPartitionsInvalidated(affectedMvPartitions.get());
+            if (marked.isEmpty()) {
+                LOG.debug("No MV partition holds the changed base partitions, 
mv={}, baseTable={}, "
+                        + "changedPartitions={}", name, baseTableInfo, 
changedPartitions);
+                return false;
             }
-            schemaChangeVersion++;
-            editLogItem = submitIvmInfoChange();
+            editLogItem = submitPartitionStatesChange(marked);

Review Comment:
   Right, and the shape is worse than the size: this is the base table's DDL 
path, under the MV write lock, so one `TRUNCATE` of a single base partition 
copies and journals every entry of an MV whose states the change did not touch. 
The map was also the caller's own by reference, so what the record carried was 
the live field.
   
   `42ed8c44575` writes the marked entries as the delta they are -- detached 
copies, the shape `raiseRebuildRequirement` already uses -- with the snapshot 
removals riding in the same record, which the replay merges over the states 
other records wrote. Alignment and the whole-MV mark keep the map, since those 
are the changes that do move every entry.
   
   
`IvmBaselineRebuildTest#testAPartitionScopedInvalidationJournalsOnlyWhatItMarked`
 captures the payload of a real `DROP PARTITION` (the edit log spied, 
delegating to the real one), asserts it carries exactly the marked partition 
with `merge` set and the removal, and then replays it over a state map that 
another record has moved in between -- so a payload that replaced the map 
instead of merging fails.
   



##########
fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java:
##########
@@ -663,54 +746,299 @@ public void alterIvmInfo(IvmInfo ivmInfo) {
      * into the journal, and a replay that replaces the field would leave the 
caller's reference
      * pointing at state that is no longer the MV's. Changing the states is 
the MV's own job, under its
      * write lock.
-     *
-     * <p>A missing map -- an image written before the field existed, or a 
non-IVM MV -- reads as empty.
      */
     public Map<String, MTMVPartitionState> getPartitionStates() {
         readMvLock();
         try {
-            if (partitionStates == null) {
-                return Collections.emptyMap();
-            }
             return 
Collections.unmodifiableMap(MTMVPartitionState.copyOf(partitionStates));
         } finally {
             readMvUnlock();
         }
     }
 
+    /**
+     * The partitions whose requirement has been raised and not met, which a 
refresh has to rebuild rather
+     * than catch up.
+     *
+     * <p>Detached names rather than the states themselves: a caller that only 
routes by them has no
+     * business holding the map the MV journals, and it needs nothing else 
from an entry.
+     */
+    public Set<String> getPartitionsNeedingRebuild() {
+        // Built before the lock, like the map getLatestEpochs returns: which 
entries go in is what needs
+        // the lock, not having somewhere to put them.
+        Set<String> res = Sets.newLinkedHashSet();
+        readMvLock();
+        try {
+            for (Entry<String, MTMVPartitionState> entry : 
partitionStates.entrySet()) {
+                if (entry.getValue().isDirty()) {
+                    res.add(entry.getKey());
+                }
+            }
+            return res;
+        } finally {
+            readMvUnlock();
+        }
+    }
+
+    /**
+     * Whether every partition the MV holds needs a rebuild, which is when a 
whole-MV refresh does nothing
+     * the per-partition routing would not.
+     *
+     * <p>A partition that holds data and does not need one makes this false: 
a whole-MV refresh would
+     * recompute it for nothing, which is the waste the per-partition routing 
exists to avoid. A partition
+     * that was never refreshed does not count against it -- a whole-MV 
refresh fills it, which its
+     * per-partition branch would do as well -- and it needs no clause of its 
own: an aligned entry is
+     * {@code {0, 1}}, so it is behind its requirement already. An MV with no 
partitions is not an
+     * escalation either.
+     *
+     * <p>Read in place rather than through {@link #getPartitionStates()}: the 
caller asks a yes/no
+     * question, and copying the map to answer it would allocate a state 
object per partition, under this
+     * lock, on every refresh -- including the ones that escalate nothing.
+     */
+    public boolean allPartitionsNeedRebuild() {
+        readMvLock();
+        try {
+            return !partitionStates.isEmpty()
+                    && 
partitionStates.values().stream().allMatch(MTMVPartitionState::isDirty);
+        } finally {
+            readMvUnlock();
+        }
+    }
+
     // ALTER_PARTITION_STATES replay applies a detached snapshot here, 
mirroring alterIvmInfo(). Live
     // invalidation changes submit their journal from the mutating method 
instead.
     //
     // A payload without the member carries no state at all, which is not the 
same as an empty map that
     // says the states are now empty: leaving them alone is the only answer 
that cannot lose state.
     public void alterPartitionStates(Map<String, MTMVPartitionState> 
partitionStates) {
-        if (partitionStates == null) {
+        replayAlterPartitionStates(partitionStates, null, false);
+    }
+
+    /**
+     * ALTER_PARTITION_STATES replay: applies the states the payload carries, 
and drops the snapshots it
+     * names. Both in one lock acquisition, because a reader that saw the new 
requirement while the
+     * snapshot was still there could let a transparent rewrite serve rows the 
rebuild has to replace.
+     *
+     * <p>A payload without the states carries none, which is not the same as 
an empty map that says the
+     * states are now empty: leaving them alone is the only answer that cannot 
lose state.
+     *
+     * <p>{@code merge} says what the payload's states are. False, which is 
what a payload written before the
+     * member existed means and what the invalidation channel still writes, 
carries the map itself and
+     * replaces. True carries only the partitions a change touched -- see 
submitPartitionStatesDelta -- so the
+     * entries it does not name belong to other records (an invalidation that 
ran during the refresh, an entry
+     * alignment added) and are merged over rather than dropped.
+     */
+    public void replayAlterPartitionStates(Map<String, MTMVPartitionState> 
partitionStates,
+            Set<String> removedSnapshotPartitions, boolean merge) {
+        writeMvLock();
+        try {
+            if (partitionStates != null) {
+                if (merge) {
+                    for (Entry<String, MTMVPartitionState> entry : 
partitionStates.entrySet()) {
+                        this.partitionStates.put(entry.getKey(), new 
MTMVPartitionState(entry.getValue()));
+                    }
+                } else {
+                    this.partitionStates = 
MTMVPartitionState.copyOf(partitionStates);
+                }
+            }
+            refreshSnapshot.removeSnapshots(removedSnapshotPartitions);
+        } finally {
+            writeMvUnlock();
+        }
+    }
+
+    /**
+     * The {@code latestEpoch} of the given MV partitions, taken under the MV 
read lock.
+     *
+     * <p>This is the value a refresh has to remember: what it read from the 
base tables is described by
+     * the requirement in force when it started reading, so writing that value 
back as the new
+     * {@code refreshEpoch} is what keeps an invalidation arriving mid-refresh 
from being swallowed. A
+     * partition without an entry is left out -- a caller writes an epoch only 
for what it captured.
+     */
+    public Map<String, Long> getLatestEpochs(Set<String> partitionNames) {
+        if (CollectionUtils.isEmpty(partitionNames)) {
+            return Collections.emptyMap();
+        }
+        // Sized before the lock: the state map is what needs it, and building 
the map is not part of that.
+        Map<String, Long> res = 
Maps.newHashMapWithExpectedSize(partitionNames.size());
+        readMvLock();
+        try {
+            for (String partitionName : partitionNames) {
+                MTMVPartitionState state = partitionStates.get(partitionName);
+                if (state != null) {
+                    res.put(partitionName, state.getLatestEpoch());
+                }
+            }
+            return res;
+        } finally {
+            readMvUnlock();
+        }
+    }
+
+    /**
+     * Brings the partition states in line with the MV's partitions: every 
partition gets an entry, and
+     * every entry whose partition is gone is dropped.
+     *
+     * <p>Alignment is what makes "the partition exists" and "the entry 
exists" the same thing, and it is
+     * why an invalidation cannot miss: rows are only written by a refresh, 
and every refresh aligns
+     * before it reads a base table, so a partition that holds rows always has 
an entry for the mark to
+     * land on. The other direction is what makes the criterion safe -- an 
entry created here describes a
+     * partition with no rows yet, so requiring one generation of it discards 
no requirement that was
+     * made earlier.
+     *
+     * <p>What it changes is journaled, because the entry has to be on disk 
before the rows it describes
+     * can be: a crash between this and the task result would otherwise leave 
a partition that holds rows
+     * with no entry at all, and every later invalidation of it would find 
nothing to land on. That is the
+     * one shape in which the criterion cannot be read -- "no entry" is 
supposed to mean "no rows" -- so
+     * the entry is made durable before any base table is read rather than 
derived again on the next run.
+     *
+     * <p>It is deliberately not a hook on every path that creates or drops a 
partition. An entry is
+     * derived state, and rebuilding it from the live partition set also 
repairs whatever a crash left
+     * behind: the drop of a partition and the removal of its entry are two 
journal records, and only
+     * their order -- partition first -- is safe, which leaves at most a stale 
entry that the next
+     * alignment drops.
+     *
+     * <p>Only an IVM MV is aligned. For a non-IVM MV the map stays as it is, 
and every reader treats
+     * "empty" and "no state" the same.
+     */
+    public void alignPartitionStates() {
+        if (!isIvm()) {
             return;
         }
+        EditLogItem editLogItem = null;
         writeMvLock();
         try {
-            this.partitionStates = MTMVPartitionState.copyOf(partitionStates);
+            // Read here rather than handed in by the caller: a caller has to 
read the names before it takes
+            // this lock, and a partition created in between -- by a 
concurrent refresh's partition sync --
+            // would then be dropped by the retainAll below, taking with it 
the state a following
+            // invalidation has to land on. The read is cheap and takes no 
lock of its own, so doing it here
+            // does not add an edge to the lock order.
+            Set<String> livePartitions = Sets.newHashSet(getPartitionNames());
+            boolean changed = 
partitionStates.keySet().retainAll(livePartitions);
+            for (String partitionName : livePartitions) {
+                if (!partitionStates.containsKey(partitionName)) {
+                    partitionStates.put(partitionName, 
MTMVPartitionState.initial());
+                    changed = true;
+                }
+            }
+            if (changed) {
+                editLogItem = 
submitPartitionStatesChange(Collections.emptySet());

Review Comment:
   Leaving this one as it stands. Alignment reads the live partition names and 
reconciles the map as a whole (`retainAll` plus a fill), so the record it 
writes is the map that reconciliation produced -- and the removal half has to 
be expressible, which a delta cannot say: an entry a merge payload omits is a 
no-op, not a drop. Splitting the two halves to make the add-only case a delta 
puts a branch in the reconciliation for an event that only fires when a 
partition is created or dropped, while the invalidation path you flagged above 
is the one that ran on every base partition DDL -- which is why that one is 
fixed.
   



##########
fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java:
##########
@@ -663,54 +746,299 @@ public void alterIvmInfo(IvmInfo ivmInfo) {
      * into the journal, and a replay that replaces the field would leave the 
caller's reference
      * pointing at state that is no longer the MV's. Changing the states is 
the MV's own job, under its
      * write lock.
-     *
-     * <p>A missing map -- an image written before the field existed, or a 
non-IVM MV -- reads as empty.
      */
     public Map<String, MTMVPartitionState> getPartitionStates() {
         readMvLock();
         try {
-            if (partitionStates == null) {
-                return Collections.emptyMap();
-            }
             return 
Collections.unmodifiableMap(MTMVPartitionState.copyOf(partitionStates));
         } finally {
             readMvUnlock();
         }
     }
 
+    /**
+     * The partitions whose requirement has been raised and not met, which a 
refresh has to rebuild rather
+     * than catch up.
+     *
+     * <p>Detached names rather than the states themselves: a caller that only 
routes by them has no
+     * business holding the map the MV journals, and it needs nothing else 
from an entry.
+     */
+    public Set<String> getPartitionsNeedingRebuild() {
+        // Built before the lock, like the map getLatestEpochs returns: which 
entries go in is what needs
+        // the lock, not having somewhere to put them.
+        Set<String> res = Sets.newLinkedHashSet();
+        readMvLock();
+        try {
+            for (Entry<String, MTMVPartitionState> entry : 
partitionStates.entrySet()) {
+                if (entry.getValue().isDirty()) {
+                    res.add(entry.getKey());
+                }
+            }
+            return res;
+        } finally {
+            readMvUnlock();
+        }
+    }
+
+    /**
+     * Whether every partition the MV holds needs a rebuild, which is when a 
whole-MV refresh does nothing
+     * the per-partition routing would not.
+     *
+     * <p>A partition that holds data and does not need one makes this false: 
a whole-MV refresh would
+     * recompute it for nothing, which is the waste the per-partition routing 
exists to avoid. A partition
+     * that was never refreshed does not count against it -- a whole-MV 
refresh fills it, which its
+     * per-partition branch would do as well -- and it needs no clause of its 
own: an aligned entry is
+     * {@code {0, 1}}, so it is behind its requirement already. An MV with no 
partitions is not an
+     * escalation either.
+     *
+     * <p>Read in place rather than through {@link #getPartitionStates()}: the 
caller asks a yes/no
+     * question, and copying the map to answer it would allocate a state 
object per partition, under this
+     * lock, on every refresh -- including the ones that escalate nothing.
+     */
+    public boolean allPartitionsNeedRebuild() {
+        readMvLock();
+        try {
+            return !partitionStates.isEmpty()
+                    && 
partitionStates.values().stream().allMatch(MTMVPartitionState::isDirty);
+        } finally {
+            readMvUnlock();
+        }
+    }
+
     // ALTER_PARTITION_STATES replay applies a detached snapshot here, 
mirroring alterIvmInfo(). Live
     // invalidation changes submit their journal from the mutating method 
instead.
     //
     // A payload without the member carries no state at all, which is not the 
same as an empty map that
     // says the states are now empty: leaving them alone is the only answer 
that cannot lose state.
     public void alterPartitionStates(Map<String, MTMVPartitionState> 
partitionStates) {
-        if (partitionStates == null) {
+        replayAlterPartitionStates(partitionStates, null, false);
+    }
+
+    /**
+     * ALTER_PARTITION_STATES replay: applies the states the payload carries, 
and drops the snapshots it
+     * names. Both in one lock acquisition, because a reader that saw the new 
requirement while the
+     * snapshot was still there could let a transparent rewrite serve rows the 
rebuild has to replace.
+     *
+     * <p>A payload without the states carries none, which is not the same as 
an empty map that says the
+     * states are now empty: leaving them alone is the only answer that cannot 
lose state.
+     *
+     * <p>{@code merge} says what the payload's states are. False, which is 
what a payload written before the
+     * member existed means and what the invalidation channel still writes, 
carries the map itself and
+     * replaces. True carries only the partitions a change touched -- see 
submitPartitionStatesDelta -- so the
+     * entries it does not name belong to other records (an invalidation that 
ran during the refresh, an entry
+     * alignment added) and are merged over rather than dropped.
+     */
+    public void replayAlterPartitionStates(Map<String, MTMVPartitionState> 
partitionStates,
+            Set<String> removedSnapshotPartitions, boolean merge) {
+        writeMvLock();
+        try {
+            if (partitionStates != null) {
+                if (merge) {
+                    for (Entry<String, MTMVPartitionState> entry : 
partitionStates.entrySet()) {
+                        this.partitionStates.put(entry.getKey(), new 
MTMVPartitionState(entry.getValue()));
+                    }
+                } else {
+                    this.partitionStates = 
MTMVPartitionState.copyOf(partitionStates);
+                }
+            }
+            refreshSnapshot.removeSnapshots(removedSnapshotPartitions);
+        } finally {
+            writeMvUnlock();
+        }
+    }
+
+    /**
+     * The {@code latestEpoch} of the given MV partitions, taken under the MV 
read lock.
+     *
+     * <p>This is the value a refresh has to remember: what it read from the 
base tables is described by
+     * the requirement in force when it started reading, so writing that value 
back as the new
+     * {@code refreshEpoch} is what keeps an invalidation arriving mid-refresh 
from being swallowed. A
+     * partition without an entry is left out -- a caller writes an epoch only 
for what it captured.
+     */
+    public Map<String, Long> getLatestEpochs(Set<String> partitionNames) {
+        if (CollectionUtils.isEmpty(partitionNames)) {
+            return Collections.emptyMap();
+        }
+        // Sized before the lock: the state map is what needs it, and building 
the map is not part of that.
+        Map<String, Long> res = 
Maps.newHashMapWithExpectedSize(partitionNames.size());
+        readMvLock();
+        try {
+            for (String partitionName : partitionNames) {
+                MTMVPartitionState state = partitionStates.get(partitionName);
+                if (state != null) {
+                    res.put(partitionName, state.getLatestEpoch());
+                }
+            }
+            return res;
+        } finally {
+            readMvUnlock();
+        }
+    }
+
+    /**
+     * Brings the partition states in line with the MV's partitions: every 
partition gets an entry, and
+     * every entry whose partition is gone is dropped.
+     *
+     * <p>Alignment is what makes "the partition exists" and "the entry 
exists" the same thing, and it is
+     * why an invalidation cannot miss: rows are only written by a refresh, 
and every refresh aligns
+     * before it reads a base table, so a partition that holds rows always has 
an entry for the mark to
+     * land on. The other direction is what makes the criterion safe -- an 
entry created here describes a
+     * partition with no rows yet, so requiring one generation of it discards 
no requirement that was
+     * made earlier.
+     *
+     * <p>What it changes is journaled, because the entry has to be on disk 
before the rows it describes
+     * can be: a crash between this and the task result would otherwise leave 
a partition that holds rows
+     * with no entry at all, and every later invalidation of it would find 
nothing to land on. That is the
+     * one shape in which the criterion cannot be read -- "no entry" is 
supposed to mean "no rows" -- so
+     * the entry is made durable before any base table is read rather than 
derived again on the next run.
+     *
+     * <p>It is deliberately not a hook on every path that creates or drops a 
partition. An entry is
+     * derived state, and rebuilding it from the live partition set also 
repairs whatever a crash left
+     * behind: the drop of a partition and the removal of its entry are two 
journal records, and only
+     * their order -- partition first -- is safe, which leaves at most a stale 
entry that the next
+     * alignment drops.
+     *
+     * <p>Only an IVM MV is aligned. For a non-IVM MV the map stays as it is, 
and every reader treats
+     * "empty" and "no state" the same.
+     */
+    public void alignPartitionStates() {
+        if (!isIvm()) {
             return;
         }
+        EditLogItem editLogItem = null;
         writeMvLock();
         try {
-            this.partitionStates = MTMVPartitionState.copyOf(partitionStates);
+            // Read here rather than handed in by the caller: a caller has to 
read the names before it takes
+            // this lock, and a partition created in between -- by a 
concurrent refresh's partition sync --
+            // would then be dropped by the retainAll below, taking with it 
the state a following
+            // invalidation has to land on. The read is cheap and takes no 
lock of its own, so doing it here
+            // does not add an edge to the lock order.
+            Set<String> livePartitions = Sets.newHashSet(getPartitionNames());
+            boolean changed = 
partitionStates.keySet().retainAll(livePartitions);
+            for (String partitionName : livePartitions) {
+                if (!partitionStates.containsKey(partitionName)) {
+                    partitionStates.put(partitionName, 
MTMVPartitionState.initial());
+                    changed = true;
+                }
+            }
+            if (changed) {
+                editLogItem = 
submitPartitionStatesChange(Collections.emptySet());
+            }
         } finally {
             writeMvUnlock();
         }
+        if (editLogItem != null) {
+            editLogItem.await();
+        }
     }
 
-    public void invalidateIvmBaseline() {
-        EditLogItem editLogItem;
+    /**
+     * The snapshots of the partitions that are clean after this result's 
epochs were applied.
+     *
+     * <p>An invalidation that reached a partition while the task ran leaves 
it dirty, and its snapshot
+     * must stay gone: dropping the entry is what keeps transparent rewrite 
away from rows the rebuild has
+     * to replace, and a result written back afterwards would undo exactly 
that. Removing only the entry
+     * keeps the rest of the map, which the removal on the invalidation side 
cannot express.
+     *
+     * <p>The caller holds the MV write lock and has already applied the 
epochs, so {@code isDirty} here
+     * reads the state the data is actually described by.
+     *
+     * <p>Only an IVM MV has partition states, so only its write-back is 
narrowed here: every entry of a
+     * non-IVM MV has no state to be dirty in and is written back as it always 
was.
+     */
+    private Map<String, MTMVRefreshPartitionSnapshot> 
snapshotsOfCleanPartitions(
+            Map<String, MTMVRefreshPartitionSnapshot> snapshots) {
+        if (MapUtils.isEmpty(snapshots)) {
+            return snapshots;

Review Comment:
   Confirmed, and it is the empty case that leaks the alias -- a non-empty map 
goes through the new map the filter builds. `after()` hands over the task's own 
`partitionSnapshots` field, which is the map the worker merges every committed 
batch into, and `executeCancelLogic` calls `after()` from the cancel thread 
(`cancel(false)`, so the execution goes on). The record is written out 
asynchronously, so what serializes is whatever the map holds by then.
   
   `42ed8c44575` closes it in both places: the empty payload is an immutable 
empty map, and `after()` hands over a copy rather than the map itself -- the 
rule `getIvmCapturedEpochs` already follows for the epochs, for the same 
reason. The copy also covers what the filter does not run on: a plain MV's 
result, and a non-empty payload the worker keeps adding to.
   
   `MTMVTest#testAnEmptySnapshotPayloadIsNotTheMapTheTaskKeepsFilling` builds a 
result from an empty map, adds an entry afterwards as the worker would, and 
asserts the journaled payload is still empty; 
`MTMVTaskTest#testTheSnapshotsHandedOverAreNotTheMapTheWorkerKeepsFilling` pins 
the handover itself.
   



-- 
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