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


##########
fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java:
##########
@@ -550,14 +676,43 @@ private PartitionRefreshPlan 
planPartitionRefresh(MTMVRefreshContext context,
                             + "does not support refreshing by partition");
         }
         try {
+            Set<String> planned = 
Sets.newLinkedHashSet(MTMVPartitionUtil.getMTMVNeedRefreshPartitions(context,
+                    relation.getBaseTablesOneLevelAndFromView()));
+            planned.addAll(rebuildRequired);
             return PartitionRefreshPlan.success(context,
-                    MTMVPartitionUtil.getMTMVNeedRefreshPartitions(context,
-                            relation.getBaseTablesOneLevelAndFromView()));
+                    excludingRebuiltPartitions(Lists.newArrayList(planned)));
         } catch (Exception e) {
             return PartitionRefreshPlan.fallback(e.getMessage());
         }
     }
 
+    /**
+     * Takes the partitions this task has already replaced out of a planned 
set.
+     *
+     * <p>The plan is computed from the snapshot the MV holds, which this task 
has not published yet, so a
+     * partition this task has just filled still looks unsynced to it. 
Refreshing it here would replace the
+     * rows it holds with a second read of the same base table, and would do 
it in the one case that has
+     * already paid for it: the same refresh falling back out of the 
incremental attempt. What the
+     * accumulator holds at this point is exactly those partitions -- the 
fallback runs after an attempt
+     * that committed nothing, and a plan is built before anything is written.
+     *
+     * <p>An explicit partition list is left alone. It is the request itself 
rather than an inference from
+     * the MV's snapshot, and the rebuild phase does not take partitions out 
of it either: what the request
+     * names is refreshed, and at worst it is refreshed twice within one task.
+     */
+    private List<String> excludingRebuiltPartitions(List<String> 
plannedPartitions) {
+        if (partitionSnapshots.isEmpty()) {
+            return plannedPartitions;
+        }
+        List<String> remaining = 
Lists.newArrayListWithCapacity(plannedPartitions.size());
+        for (String partitionName : plannedPartitions) {
+            if (!partitionSnapshots.containsKey(partitionName)) {

Review Comment:
   [P1] Preserve a current non-PCT baseline when excluding rebuilt partitions 
from fallback. Suppose a partitioned F JOIN D aggregate has p1/p2 and D changes 
10 to 20 while p1 is dirty. The pre-rebuild of p1 reads D SNAPSHOT at the old 
offset (10). If the IVM delta returns a non-COMPLETE fallback reason, this 
helper removes p1 from the PARTITIONS plan, leaving only p2; that partial 
overwrite also reads D SNAPSHOT at 10, so the task can report SUCCESS with both 
partitions at 10 while recording D's current version (20) for them. Later sync 
checks can then skip the still-pending D delta and transparent rewrite can 
serve wrong rows. The old full-scope fallback reset D and repaired both 
partitions; the earlier :815 thread covered its repeated overwrite, not this 
new stale-success outcome. Reconcile D's pending stream change or take a full 
current-snapshot rebuild before publishing success, and test exact rows after a 
forced non-COMPLETE IVM fallback.



##########
fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java:
##########
@@ -947,67 +1297,118 @@ private MTMVRelatedTableIf findPctTable(BaseTableInfo 
baseTableInfo) {
         return null;
     }
 
+    private EditLogItem submitIvmInfoChange() {
+        // The caller has already mutated ivmInfo under the MV write lock. 
Submit its snapshot directly;
+        // replay later applies the payload through alterIvmInfo().
+        AlterMTMV alterMTMV = new AlterMTMV(
+                new TableNameInfo(getQualifiedDbName(), getName()), 
MTMVAlterOpType.ALTER_IVM_INFO);
+        alterMTMV.setIvmInfo(ivmInfo);
+        return submitAlterLog(alterMTMV);
+    }
+
     /**
-     * Release the IVM baseline barrier after the partitions it named have 
been rebuilt, or after
-     * partition sync removed them (a dropped partition resolves its own 
entry: the partition and its
-     * IVM offsets are both gone).
+     * Raises the requirement of the given MV partitions, so the next refresh 
rebuilds them.
      *
-     * <p>Guarded by schemaChangeVersion, like {@link 
#persistIvmBaselineGuard}: a base-table change
-     * landing while the rebuild runs carries its own barrier entry, and a 
blind clear would swallow
-     * it. Failing instead preserves that entry -- the next refresh rebuilds 
it together with the
-     * partitions this task handled.
-     *
-     * <p>Journals the new state right away, like every other ivmInfo mutation 
here. A task that dies
-     * before {@link #addTaskResult} would otherwise leave the release in 
memory only, and a restart
-     * would resurrect the barrier from disk.
+     * <p>This is an invalidation-shaped mutation, journaled as the whole 
state map before whatever needs
+     * it is done. A caller about to make a partition's rows unusable says so 
with it: the raised
+     * requirement survives a crash, so a refresh that never got to publish 
its rebuild leaves partitions
+     * naming a generation they do not hold, and the next refresh rebuilds 
them.
      */
-    public void releaseIvmBaselineRebuild(long expectedSchemaChangeVersion) 
throws JobException {
+    public void markPartitionsForRebuild(Set<String> partitionNames) {
+        if (CollectionUtils.isEmpty(partitionNames)) {
+            return;
+        }
         EditLogItem editLogItem;
         writeMvLock();
         try {
-            if (ivmInfo == null || !ivmInfo.isBaselineRebuildRequired()) {
-                // Nothing to release: skip both the mutation and the journal 
entry. Any base-table
-                // change that raced us in is still caught by 
validateIvmRefreshStart() below.
-                return;
+            boolean changed = false;
+            for (String partitionName : partitionNames) {
+                MTMVPartitionState state = partitionStates.get(partitionName);
+                if (state == null) {
+                    // A partition dropped since the caller planned it has no 
rows to protect.
+                    continue;
+                }
+                state.setLatestEpoch(state.getLatestEpoch() + 1);
+                changed = true;
             }
-            if (schemaChangeVersion != expectedSchemaChangeVersion) {
-                throw new JobException("Base table metadata changed before IVM 
baseline refresh, mv="
-                        + getName());
+            if (!changed) {
+                return;
             }
-            ivmInfo.clearBaselineRebuild();
-            editLogItem = submitIvmInfoChange();
+            editLogItem = submitPartitionStatesChange(Collections.emptySet());
         } finally {
             writeMvUnlock();
         }
         editLogItem.await();
     }
 
-    public void persistIvmBaselineGuard(RefreshMode refreshMode, Set<String> 
baselinePartitions,
-            long expectedSchemaChangeVersion) throws JobException {
+    /**
+     * Raises the requirement of the given MV partitions that do not name one, 
and reports what each of those
+     * partitions now names.
+     *
+     * <p>This is what a refresh about to replace a partition says about it, 
and it has to be said before that
+     * replacement reads anything: an overwrite is two halves -- the rows are 
committed into temporary
+     * partitions, and a swap publishes them -- so a refresh that dies in 
between leaves the live partition
+     * holding the rows it had while whatever its read consumed, the offsets 
of the streams it read among
+     * them, is already committed with the first half. The epochs a refresh 
records ride with its result, and
+     * a refresh that never returns records none, so without this nothing 
would say the partition owes the
+     * rebuild and the next refresh would read on from an offset past a change 
the partition never received.
+     *
+     * <p>What the caller gets back is the requirement it raised, which is the 
ceiling its write-back is
+     * clamped to; see MTMVTask's captured epochs. A partition that already 
names a requirement is left alone
+     * and is not part of that result: it names the requirement the refresh 
answers for, and the caller must
+     * record what it read rather than what it found. Raising it again would 
move it above that, and the
+     * partition would be rebuilt a second time for nothing. The record it 
submits still carries the whole
+     * map -- that is what this channel carries -- but it is submitted only 
when something was raised.
+     *
+     * <p>This differs from {@link #markPartitionsForRebuild} on purpose: that 
one is an invalidation, and it
+     * raises the requirement of every partition it names because it has to 
outrank a refresh already
+     * running. This one is a refresh's own record of what it is about to do, 
and a partition that already
+     * names such a requirement does not need a second one.
+     *
+     * <p>Which partitions need it is decided under the same lock as the 
raise. Read outside it, a mark
+     * landing in between would leave this call blind to a requirement it then 
raises above, and the caller
+     * would record its own value as met for a change that arrived after it 
read.
+     */
+    public Map<String, Long> raiseRebuildRequirement(Set<String> 
partitionNames) {
+        if (CollectionUtils.isEmpty(partitionNames)) {
+            return Collections.emptyMap();
+        }
+        Map<String, Long> raised = 
Maps.newHashMapWithExpectedSize(partitionNames.size());
         EditLogItem editLogItem;
         writeMvLock();
         try {
-            if (schemaChangeVersion != expectedSchemaChangeVersion) {
-                throw new JobException("Base table metadata changed before IVM 
baseline refresh, mv=" + getName());
+            for (String partitionName : partitionNames) {
+                MTMVPartitionState state = partitionStates.get(partitionName);
+                if (state == null || state.isDirty()) {
+                    // Dropped since the caller planned it, or already naming 
a requirement of its own.
+                    continue;
+                }
+                state.setLatestEpoch(state.getLatestEpoch() + 1);
+                raised.put(partitionName, state.getLatestEpoch());
             }
-            if (refreshMode == RefreshMode.COMPLETE) {
-                ivmInfo.requireCompleteBaselineRebuild();
-            } else {
-                
ivmInfo.addPendingBaselineRebuildPartitions(baselinePartitions);
+            if (raised.isEmpty()) {
+                return Collections.emptyMap();
             }
-            editLogItem = submitIvmInfoChange();
+            editLogItem = submitPartitionStatesChange(Collections.emptySet());

Review Comment:
   [P2] Persist only the selected partition's new rebuild requirement here. 
Each IVM PARTITIONS refresh of one clean partition calls this method before 
overwrite, then submitPartitionStatesChange copies every partition state under 
mvRwLock and serializes the whole map to the edit log. On an MV with many 
partitions, refreshing just one becomes a large lock-held copy and journal 
record each time this path runs. The earlier ADD_TASK thread addressed that 
later task-result payload; this new routine pre-overwrite record has the same 
scaling problem. Use a merge-on-replay epoch delta for the selected entries and 
cover one-partition journal size on a large MV.



##########
fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java:
##########
@@ -715,12 +768,47 @@ 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);

Review Comment:
   [P1] Do not publish this rebuild as fresh if the following delta fails. A 
strict INCREMENTAL can rebuild dirty p1 using a non-PCT table's SNAPSHOT at its 
old stream offset, while generatePartitionSnapshots records that table's 
current version. For a partitioned F JOIN D aggregate with D changing 10 to 20, 
p1's overwrite still holds 10; if the subsequent IVM delta returns a fallback 
reason, strict mode fails, but onFail publishes p1's captured epoch and current 
D snapshot. FAILED MVs remain eligible for transparent rewrite, and the 
snapshot comparisons then admit p1 although the base query returns 20. The old 
strict path rejected this pending rebuild. Keep such partitions 
dirty/unavailable until the required delta commits, or tag the rebuild with the 
historical source snapshot; add a forced-delta-failure regression that checks 
rewrite eligibility and rows.



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