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


##########
fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java:
##########
@@ -928,8 +1217,26 @@ private void 
executePartitionBasedRefresh(MTMVRefreshContext context, RefreshMod
                         mtmv.getName(), getTaskId(), e);
                 throw new JobException(e.getMessage(), e);
             }
-            completedPartitions.addAll(execPartitionNames);
-            partitionSnapshots.putAll(execPartitionSnapshots);
+            recordRefreshCompleted(execPartitionNames);
+            // What this batch may record: everything it replaced, unless the 
read behind it answered with a
+            // base table the MV does not partition by as of an older state 
than that table is in now -- which
+            // a partial read does for the partitions holding data the stream 
offset has not consumed, and it
+            // is the delta that follows that brings the table up to date. For 
the partitions this batch
+            // replaced, that delta does not apply: they are kept out of its 
scope, so a delta computed against
+            // the rows a rebuild replaced cannot be applied twice. Recording 
them would say they hold the
+            // table's current state when they hold that image, and nothing 
would plan them again. Left
+            // unrecorded, the requirement stays, the MV keeps the snapshot it 
has rather than one claiming the
+            // current state, and a later refresh rebuilds them -- by then the 
offset has been consumed, so a
+            // rebuild reads the table as it is. What the read answered with 
is recorded where the read is
+            // planned; see IvmFullRefreshMTMV#readsAnOlderImage.
+            if 
(rewriteContext.map(IvmRewriteContext::isReadFromAStreamOffset).orElse(false)) {

Review Comment:
   Right, and the evidence is in the suite my last commit added -- I read `kept 
out of its scope` as `the delta does not write to it`, when the scope only 
decides what the task records. The delta's target is the partitions the change 
touches, and that includes the ones the rebuild replaced; my own case showed 
the rebuilt partition moving from `10` to `20` during the successful delta 
while I was still withholding its epoch and snapshot. Holding them was 
therefore a full rebuild of those partitions on every refresh, for as long as a 
table they read keeps changing -- which is what you describe.
   
   `635c078bcb4` holds them instead of dropping them, and publishes them when 
the incremental attempt's delta succeeds. Two things I checked before doing 
that, since it is the direction that can publish too much:
   
   * the reason a batch is held is always a table the MV does not partition by 
-- the flag is set only where a snapshot read is bound, and the table the MV 
partitions by is read with `RESET` -- so the delta that redeems it is that 
table's delta, over the partitions the change maps to;
   * a held partition the change does not touch is one whose rows the change 
does not reach, so publishing it is not claiming rows it does not have.
   
   When the delta does not run or fails, the records stay held, which is the 
case the holding exists for. The visible effect in the suite added with the 
last commit: its second pass now reports no scope and no rebuild where it used 
to report one partition and one, so the recovery is finished there rather than 
deferred by one -- the golden is regenerated and reviewed line by line for that.
   



##########
fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java:
##########
@@ -792,13 +912,159 @@ private IvmIncrRefreshResult 
executeSingleIvmAttempt(MTMVRefreshContext refreshC
         }
         if (ivmResult.isSuccess()) {
             this.partitionSnapshots.putAll(capturedSnapshots);
-            this.completedPartitions.addAll(needRefreshPartitions);
+            recordRefreshCompleted(incrementalScope);
+            commitCapturedEpochs(capturedEpochs);
             LOG.info("IVM incremental refresh succeeded for mv={}, taskId={}",
                     mtmv.getName(), getTaskId());
         }
         return ivmResult;
     }
 
+    /**
+     * Adds a phase's scope to what this task reports as refreshed, and the 
partitions it committed to what
+     * this task reports as done. Both are the task's; see the fields for why.
+     */
+    private void recordRefreshScope(Collection<String> partitions) {
+        // Created by the first phase that records one: a task that has not 
refreshed anything reports
+        // nothing, which is the same state as one that has not run yet. 
Concurrent because the columns it
+        // feeds are read while the worker fills them -- the tasks() table 
function reports a running task
+        // -- and ordered because what this reports is persisted and has to 
read the same on every look.
+        if (needRefreshPartitions == null) {
+            needRefreshPartitions = new ConcurrentSkipListSet<>();
+        }
+        needRefreshPartitions.addAll(partitions);
+    }
+
+    private void recordRefreshCompleted(Collection<String> partitions) {
+        if (completedPartitions == null) {
+            completedPartitions = new ConcurrentSkipListSet<>();
+        }
+        completedPartitions.addAll(partitions);
+    }
+
+    /**
+     * Says durably that the parts of the scope that do not name a rebuild 
requirement have to be rebuilt, and
+     * brings the epochs this phase records in line with it.
+     *
+     * <p>Raised here rather than by each caller of the executor, like the 
scope above: this is the phase that
+     * replaces partitions, and a caller that forgot would leave a partition 
whose rows were never published
+     * looking caught up. The partitions a caller has already made dirty are 
left as they are, which is what
+     * makes raising it here harmless for them: a whole-MV attempt marks its 
scope before it reconciles the
+     * streams, and the incremental attempt rebuilds the partitions an 
invalidation marked.
+     * See MTMV#raiseRebuildRequirement.
+     *
+     * <p>What that call reports is what a partition it raised now names, and 
this phase is clamped to it. The
+     * clamp cannot stay at what an earlier attempt planned: that value sits 
below the requirement this phase
+     * has just raised, so the epochs recorded here would leave the partition 
dirty after it was replaced, and
+     * every refresh after it would rebuild the same partitions again.
+     *
+     * <p>A partition that already named a requirement keeps the entry it has, 
and is deliberately not moved
+     * up to what it names now. The entry is the value the routing decision 
saw, and it is what keeps a mark
+     * landing between that decision and this phase's read from being recorded 
as met by a replacement that
+     * read before the change it made. Leaving it where it is costs one 
rebuild; moving it up could cost the
+     * change.
+     */
+    private void raiseRequirementForRefreshScope(Collection<String> 
partitions) {
+        if (!mtmv.isIvm()) {
+            // A plain MV has no streams to read and its epochs record 
nothing: the sync criterion plans it
+            // again on its own until its snapshots are published.
+            return;
+        }
+        
ivmPlannedEpochs.putAll(mtmv.raiseRebuildRequirement(Sets.newHashSet(partitions)));
+    }
+
+    /**
+     * Brings the routing decision up to date after a retry has synchronized 
and aligned the MV's partitions.
+     *
+     * <p>A partition the alignment creates is dirty by construction -- {@code 
{0, 1}}, behind its
+     * requirement -- and the decision was taken before it existed. Reading 
the states again is what makes
+     * the retried attempt treat it as such: it joins the dirty set, so the 
incremental attempt leaves it
+     * out rather than recording what a delta captured as the partition being 
caught up, and it gets the
+     * entry the batches are clamped against, so a mark landing later in this 
task cannot be written back as
+     * satisfied either.
+     *
+     * <p>A partition this leaves dirty is not rebuilt here. The rebuild phase 
has already run, and the
+     * partition the alignment created holds no rows yet -- there is nothing 
to replace in it -- so what it
+     * needs is a build, which is what the next refresh's rebuild gives it.
+     *
+     * <p>The planned value of a partition that already had one is kept: that 
is the value the routing
+     * decision was made on, which is what the clamp is for.
+     */
+    private void adoptPartitionsCreatedByTheRetry(Set<String> dirtyPartitions) 
{
+        Set<String> livePartitionNames = mtmv.getPartitionNames();
+        for (Entry<String, MTMVPartitionState> entry : 
mtmv.getPartitionStates().entrySet()) {
+            if (!livePartitionNames.contains(entry.getKey())) {
+                continue;
+            }
+            ivmPlannedEpochs.putIfAbsent(entry.getKey(), 
entry.getValue().getLatestEpoch());
+            if (entry.getValue().isDirty()) {
+                dirtyPartitions.add(entry.getKey());
+            }
+            // A partition the alignment has just created -- the state an 
entry starts with,
+            // MTMVPartitionState.initial() -- is a
+            // partition of this name that the retry's partition sync 
recreated: the partition the captures and
+            // snapshots this task holds for that name describe is gone, and 
the one that took its name holds
+            // nothing. Writing those back would credit the new one with what 
the old one held, which is worse
+            // than a wrong number: a partition clean at an epoch a later 
change only raises to is one no
+            // refresh rebuilds, so the rows the recreation removed would be 
published as current.
+            if (entry.getValue().getRefreshEpoch() == 0 && 
entry.getValue().getLatestEpoch() == 1) {

Review Comment:
   Right, and the shape of the state is exactly what cannot tell the two apart: 
a partition this task rebuilt has its rows and its capture committed, and its 
state still reads the one an entry starts with until the task result writes the 
epochs back -- so my fence took the partition I had just repaired for one that 
had been recreated, and dropped what it had committed. Leaving a repaired 
partition unrecorded is not free either: it is rebuilt in full by the next 
refresh, and out of transparent rewrite until then.
   
   `635c078bcb4` fences by identity instead: `commitCapturedEpochs` records the 
id of the MV partition the capture was taken under, and 
`adoptPartitionsCreatedByTheRetry` drops the capture, snapshot and record only 
when the live partition of that name carries a different id -- which is what a 
sync drop-and-re-add leaves behind. A partition the task itself rebuilt keeps 
its id, so what it committed stands.
   
   `MTMVTaskTest` covers both directions now: same id with a 
freshly-aligned-looking state keeps everything, different id drops all three. 
The end-to-end case you describe -- a newly aligned `p`, a `q` arriving later, 
the retry syncing `q` and leaving `p` alone -- I have not built; it needs 
repeatedly inducing a `MV_PARTITION_NOT_FOUND` delta failure by dropping and 
re-adding populated base partitions, which is why I pinned the mechanism at the 
unit level. Say the word and I will add it as its own suite rather than growing 
one that already drives several shapes.
   



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