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


##########
fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java:
##########
@@ -259,7 +260,25 @@ private PartitionPlanningException(String message, 
Throwable cause) {
     // Written by the executing (Disruptor worker) thread via the 
executeCommand consumer
     // callback and read by the cancel (command) thread, so it must be 
volatile.
     private volatile StmtExecutor executor;
-    private Map<String, MTMVRefreshPartitionSnapshot> partitionSnapshots;
+    // What this task has committed, per MV partition: the snapshot each 
partition's rows were read at.
+    // One accumulator for the whole task rather than one per phase, because 
that is what the MV publishes
+    // at the end of it -- a phase that started from empty would publish its 
own work and drop the work of
+    // the phases before it, leaving partitions a preceding rebuild replaced 
looking unsynced.
+    private Map<String, MTMVRefreshPartitionSnapshot> partitionSnapshots = 
Maps.newConcurrentMap();
+    // The requirement each refreshed partition was read under, captured 
before the base tables were read
+    // and recorded only once that batch's data committed (see 
commitCapturedEpochs). In memory only: the
+    // journal carries the resulting states, and a replay applies those 
instead of recomputing anything.
+    private transient Map<String, Long> ivmCapturedEpochs = Maps.newHashMap();

Review Comment:
   One further instance of this is closed since this thread was opened: the 
snapshot map the result carries was handed over live as well, and `after()` now 
passes a copy of it -- the same rule `getIvmCapturedEpochs` already followed 
for the epochs. So what the cancel thread publishes is now a detached read of 
the committed work on every path: a copy of the epochs, a copy of the 
snapshots, and a write-back that filters both to the partitions that are still 
clean.
   
   What the finding's second half asks for -- publishing and releasing only 
after the worker quiesces -- is the framework change described above: STOP 
reaches `cancel(false)` through `AbstractJob`, so a cancelled task publishes 
before its worker stops and can see resources the cancel path has released, and 
deciding what a cancelled refresh publishes (and whether STOP waits) is a 
decision for every job type that shares the framework rather than for this PR. 
I am resolving this thread on that division: the publication side is fixed and 
pinned, the lifecycle side is tracked as its own change.
   



##########
fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java:
##########
@@ -684,33 +759,211 @@ public Map<String, MTMVPartitionState> 
getPartitionStates() {
     // 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) {
-            return;
-        }
+        replayAlterPartitionStates(partitionStates, null);
+    }
+
+    /**
+     * 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.
+     */
+    public void replayAlterPartitionStates(Map<String, MTMVPartitionState> 
partitionStates,
+            Set<String> removedSnapshotPartitions) {
         writeMvLock();
         try {
-            this.partitionStates = MTMVPartitionState.copyOf(partitionStates);
+            if (partitionStates != null) {
+                this.partitionStates = 
MTMVPartitionState.copyOf(partitionStates);
+            }
+            refreshSnapshot.removeSnapshots(removedSnapshotPartitions);
         } finally {
             writeMvUnlock();
         }
     }
 
-    public void invalidateIvmBaseline() {
-        EditLogItem editLogItem;
+    /**
+     * 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(Set<String> livePartitionNames) {
+        if (!isIvm()) {
+            return;
+        }
+        // Copied up front: callers pass what OlapTable holds, and that is 
mutated under the table's own
+        // write lock, not this one. Iterating the live collection could see 
it change.
+        Set<String> livePartitions = Sets.newHashSet(livePartitionNames);
+        EditLogItem editLogItem = null;
         writeMvLock();
         try {
-            if (ivmInfo == null) {
-                ivmInfo = new IvmInfo();
+            boolean changed = 
partitionStates.keySet().retainAll(livePartitions);

Review Comment:
   Recording the state of this one as it stands: the read of the live partition 
names was moved into `alignPartitionStates` itself, under the same MV lock as 
the map it edits, so the stale-snapshot half is gone and a caller cannot hand 
in names it read before the lock.
   
   What is left is the window between that read and the `retainAll`, both 
adjacent statements under the MV lock -- which does not close it, because the 
MV's partition map is mutated under the table's lock rather than this one. 
Closing it means removing the entry where the partition DDL happens instead of 
deriving it here, which is a change to the partition lifecycle rather than to 
this one. I am resolving the thread on that division; say the word if you want 
that change taken up separately.
   



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