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


##########
fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/NormalizeOlapTableStreamScan.java:
##########
@@ -262,6 +263,15 @@ private Plan makeSnapshotScan(LogicalOlapTableStreamScan 
scan, CascadesContext c
         OlapTable baseTable = streamWrapper.getBaseTable();
         List<Slot> originSlots = scan.getOutput();
         selectedPartitionIds = 
streamWrapper.filterConsumedPartitionIds(selectedPartitionIds);
+        // What this read is about to answer with is recorded here, where the 
plan is final and the read states
+        // are in place, and before the read is built: the refresh that reads 
it back may not record a
+        // partition it replaced through an image of the table as of the 
offset as caught up -- what it wrote
+        // is that image, and the delta that would bring the table up to date 
does not apply to it. See
+        // MTMVTask#executePartitionBasedRefresh.
+        if (streamWrapper.readsSnapshotOfAnOlderImage(selectedPartitionIds)) {

Review Comment:
   Right, and the reason it slipped is in the question it asked. It asked 
whether the read answered with the image at the offset, and took everything 
else for the table as it is now. A read that drops a partition with no 
consumption baseline answers with neither: there is no offset to read that 
partition from, and rows it holds by now are simply missing from the answer. So 
the partition was recorded as caught up while the rows it should have joined in 
were absent, and a current snapshot was published over it.
   
   `81fedd86e27` asks the question the answer deserves: 
`OlapTableStreamWrapper#answersWithTheCurrentTable`. It says no for a partition 
read behind its offset, and no for a partition the read dropped **that holds 
rows** -- the read is missing them either way. It says yes for a partition the 
read dropped that holds nothing, which is what the round before this one asked 
for: withholding there withholds a partition that was read as the table is, and 
costs a rebuild it does not owe.
   
   The two ways are pinned per key type in `OlapTableStreamSnapshotReadTest`, 
including the dropped-partition-with-rows and dropped-partition-without-rows 
pair, and the suites around these reads are green: 
`test_ivm_partial_rebuild_is_not_recorded`, 
`test_ivm_fallback_keeps_the_whole_scope`, 
`test_ivm_fallback_stream_multi_batch` and `..._dup`, `test_ivm_snapshot`, 
`test_ivm_partition_sync_limit_with_window`. The judgement stays per batch, so 
a batch that answers incompletely keeps every partition it replaced unrecorded; 
that errs towards one more rebuild and never towards rows.
   



##########
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;
+        }
+        Map<String, MTMVRefreshPartitionSnapshot> res = 
Maps.newHashMapWithExpectedSize(snapshots.size());
+        for (Entry<String, MTMVRefreshPartitionSnapshot> entry : 
snapshots.entrySet()) {
+            MTMVPartitionState state = partitionStates.get(entry.getKey());
+            // No entry means the partition was created after the alignment, 
so it can only hold rows this
+            // task wrote; a dirty one needs its rebuild before anything may 
read it through the MV.
+            if (state == null || !state.isDirty()) {
+                res.put(entry.getKey(), entry.getValue());
+            }
+        }
+        return res;
+    }
+
+    /**
+     * Records the epochs the given partitions were read at, which is how a 
refresh turns a requirement
+     * into the state of the data.
+     *
+     * <p>Only {@code refreshEpoch} is written: a refresh writes back the 
requirement it captured, and the
+     * requirement may have been raised again since that capture. A payload 
built from the captured map
+     * would overwrite the newer value and lose the rebuild it asks for, so 
{@code latestEpoch} is left
+     * alone here.
+     *
+     * <p>The caller holds the MV write lock (it is applied together with the 
rest of a task result).
+     */
+    private void applyRefreshedEpochs(Map<String, Long> capturedEpochs) {
+        if (MapUtils.isEmpty(capturedEpochs)) {
+            return;
+        }
+        for (Entry<String, Long> entry : capturedEpochs.entrySet()) {
+            MTMVPartitionState state = partitionStates.get(entry.getKey());
+            if (state == null) {
+                // The partition was dropped while the task ran, so its state 
went with it.
+                continue;
+            }
+            state.setRefreshEpoch(entry.getValue());

Review Comment:
   Right. The captures and snapshots a task holds are keyed by partition name, 
and a retry's partition sync can drop a partition and add one of the same name 
back -- which the alignment gives `{0, 1}`, empty. Writing the old captures 
back credits the new partition with what the old one held, and the failure mode 
is not only the wrong number: a partition clean at an epoch that a later change 
only raises to is one no refresh rebuilds, so the rows the recreation removed 
are published as current.
   
   `81fedd86e27` drops both, in the hook that already re-reads the partitions 
after a retry (`adoptPartitionsCreatedByTheRetry`): a live partition whose 
state is the one an entry starts with -- `MTMVPartitionState.initial()` -- is 
one the alignment just created, so the capture and the snapshot this task holds 
for that name go. The snapshot has to go with the capture: leaving it would let 
the sync judgement call the recreated, empty partition caught up even with its 
epoch no longer written back.
   
   `MTMVTaskTest` pins it: with `p1` recreated and `p2` the partition the 
rebuild replaced, `p1`'s capture and snapshot are dropped, `p2`'s stand, and 
both are dirty.
   
   What I did not build is the end-to-end case you describe -- two retries with 
a same-name re-add. It needs repeatedly inducing a "missing MV partition" delta 
failure by dropping and re-adding populated base partitions, and the mechanism 
is small enough that I pinned it at the unit level instead; if you want it as a 
suite, say so and I will add one on its own rather than growing a suite that 
already drives several shapes. The reporting side is left as it was on purpose: 
`CompletedPartitions` still names the partition this task refreshed under that 
name, which is what it did; the epochs are what must not be inherited.
   



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