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


##########
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:
   [P2] Stop withholding a rebuilt partition after its successful delta catches 
it up. `incrementalScope.removeAll(dirtyPartitions)` only changes task 
bookkeeping; the actual delta insert uses 
`generateRelatedBasePartitionIds(mtmv)` and can update those partitions. The 
new partial-rebuild regression itself shows p1 move from D.v=10 to 20 during 
the successful delta, yet this branch omits p1's epoch/snapshot and the next 
task rebuilds it again. With p2 clean and another non-PCT update before every 
refresh, p1 remains dirty and is fully overwritten every time. Track the 
physical delta target and publish a rebuilt partition once its older-image 
changes have been applied (or otherwise ensure this path converges under 
continuous updates).



##########
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:
   [P2] Clear captures only for partitions the retry actually replaced. A newly 
aligned p starts at `{0,1}`; its dirty rebuild can commit rows and capture 
epoch/snapshot while the live state remains `{0,1}` until ADD_TASK. If a 
different base partition q arrives after the initial sync and makes the delta 
return `MV_PARTITION_NOT_FOUND`, retry sync adds q but leaves p intact. This 
predicate then erases p's committed capture/snapshot, and the retried delta 
does not restore that bookkeeping, so a successful task leaves p dirty and 
unavailable for transparent rewrite until another full overwrite. Compare 
partition identity or the sync's changed-name set; cover an unchanged `{0,1}` 
partition alongside a newly added q.



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