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


##########
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:
   Confirmed, and measured rather than argued. `f4e35a1eb95` is the fix.
   
   **What happens.** The incremental attempt rebuilds the partitions the 
criterion says must be rebuilt before its delta, and that rebuild is a partial 
read: `D` is read through its stream **at the offset** -- the image as of the 
last consumption -- for the partitions holding data the offset has not 
consumed. The delta that follows is what brings `D` up to date, and it keeps 
the rebuilt partitions out of its scope on purpose. So when the delta falls 
back instead, the partition plan answers the snapshot question again, the 
partition the rebuild just replaced still looks unsynced to it, and 
`excludingRebuiltPartitions` takes it out of the plan. The plan is then 
partial, which is what changes the read for the partitions that remain: the 
same old image of `D`, with the current state recorded for every partition it 
refreshes. The task reports `SUCCESS`, every partition holds `D` as of the old 
image, every partition is recorded as current, and the partition that was taken 
out keeps what th
 e rebuild gave it from an image the refresh has now moved the offset past.
   
   **The fix.** A plan that covers the whole MV is left alone. Such a plan is 
read through `RESET`: the current rows, with the stream offsets advanced to the 
end, which only a read that saw the whole table may do. Keeping it whole is 
what keeps the fallback's read of `D` current, and the partitions the rebuild 
replaced are replaced again through that current read rather than skipped. The 
exclusion stays for a plan that is already partial, where the read it would 
give the skipped partition is the same snapshot read the rebuild itself used, 
so repeating it would redo the work for the same rows.
   
   **Cost, stated rather than hidden.** When the earlier phase had itself 
planned the whole MV, its read was the `RESET` one and the partitions it 
replaced were current, so the second overwrite is not needed there. Telling the 
two apart would take remembering, per partition, how its replacement read; the 
coarse rule's price is one overwrite, and the direction it errs in is the one 
that never leaves rows behind.
   
   **Evidence.** The suite added with it, 
`test_ivm_fallback_keeps_the_whole_scope`, drives the shape end to end -- a 
truncated base partition for the requirement, an upserted dimension for the 
pending change, the forced non-`COMPLETE` fallback reason for the fallback 
itself -- and asserts the task's scope, the rows, and that the refresh after it 
has nothing left to do. It is its own control: with the exclusion unrestored, 
the fallback's task reports `PARTIAL` over one partition instead of `COMPLETE` 
over the whole MV, and both rows keep the dimension's old value (`10` where the 
base query returns `20`) while the task reports success. Unit tests cover both 
directions of the rule in `MTMVTaskTest`, and the two suites that drive the 
same forced-fallback path (`test_ivm_fallback_stream_multi_batch`, `..._dup`) 
pass unchanged.
   



##########
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:
   Confirmed. `058c80e8c96` fixes it, and testing the judgement itself turned 
up one more thing worth recording.
   
   **What happens.** The rebuild the incremental attempt runs is a partial 
read, so `D` answers with the image at the offset whenever it holds data the 
offset has not consumed. The delta that would bring `D` up to date never 
applies to the rebuilt partitions -- they are kept out of its scope so a delta 
computed against the rows a rebuild replaced cannot be applied twice -- so 
recording them says they hold `D`'s current state while they hold that image, 
and nothing plans them again: the view keeps rows the base query does not 
return, and `SyncWithBaseTables` admits them for transparent rewrite.
   
   **The fix.** A batch whose read answered with an older image than the table 
is in now is not recorded at all -- no epoch, no snapshot. The requirement 
stays, the MV keeps the snapshot it has rather than one claiming the current 
state, and a later refresh rebuilds those partitions; by then the offset has 
been consumed, so the rebuild reads the table as it is, which is what makes the 
passes converge. A whole-MV read is not in that position: it reads those tables 
as they are now and advances their offsets, so it is recorded as before.
   
   **Where the answer comes from is the part I would look at.** It is not 
re-derived by the refresh -- an earlier revision of this asked the streams 
itself, which in cloud mode meant an RPC for the table stream read state. The 
read records it where it is planned instead: `IvmFullRefreshMTMV` decides the 
read mode, so it asks the stream scan what a snapshot read of it answers with 
(`OlapTableStreamWrapper#readsSnapshotOfAnOlderImage`, the same split 
`NormalizeOlapTableStreamScan` makes when it binds the read) and marks it on 
the `IvmRewriteContext`. A refresh owns one context per statement it runs, so 
the answer reaches the executor without a field on the `StatementContext`, 
without a planner hook, and without reading a stream state anywhere new.
   
   **What testing it turned up, and now pinned.** A duplicate-key table is read 
bounded by its offset for every partition rather than split into the two kinds 
a merge-on-write one is, so its answer is a different expression. My first 
attempt read the pairs `getHistoryPartitionOffsets` returns, whose left side is 
deliberately null -- which made the judgement **constant false**. It recorded 
nothing and failed nothing, so the suite covering that shape 
(`test_ivm_fallback_stream_multi_batch_dup`) passed for the wrong reason while 
the hole stayed open for duplicate-key MVs; a direct test of the judgement is 
what exposed it. It reads the wrapper's own offsets now, and 
`OlapTableStreamSnapshotReadTest` pins both key types, the caught-up case, and 
the partition with no consumption baseline -- counting one of those as behind 
an offset it does not have is what first turned that suite red.
   
   **Tests.** `test_ivm_partial_rebuild_is_not_recorded` drives the shape end 
to end and pins the convergence: the strict refresh whose delta fails leaves 
the rebuilt partition unrecorded; the refresh after it consumes the dimension's 
change and is still not entitled to record it (the count says so); the one 
after that reads the dimension as it is, lands the rows, and records; a fourth 
has nothing left to do. With the rule removed the second step reports nothing 
rebuilt instead of one partition.
   



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