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


##########
fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/IvmFullRefreshMTMV.java:
##########
@@ -104,6 +104,16 @@ private Plan rewriteScan(LogicalOlapScan scan, 
IvmRewriteContext rewriteContext)
 
         OlapTableStream stream = 
IvmUtil.getIvmStream(rewriteContext.getMtmv(), baseTable);
         OlapTableStreamWrapper streamWrapper = new 
OlapTableStreamWrapper(stream, baseTable, selectedPartitionIds);
+        // A snapshot read of a table the MV does not partition by answers 
with the image at the stream offset
+        // for the partitions behind it, and with the table as it is for the 
rest; which of the two each
+        // partition gets is decided here, so it is recorded here. The refresh 
that reads this back may not
+        // record a partition it replaced through the first kind 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 (readMode == StreamReadMode.SNAPSHOT
+                && 
streamWrapper.readsSnapshotOfAnOlderImage(selectedPartitionIds)) {

Review Comment:
   [P1] Install cloud read states before classifying the snapshot. This rule 
only queues `CloudTableStreamReadStateHook`, then calls 
`readsSnapshotOfAnOlderImage`; `NereidsPlanner.afterAnalyze` installs the 
authoritative states after analyzer rules finish. In cloud mode the wrapper 
therefore falls back to stale local offsets, while 
`NormalizeOlapTableStreamScan` later uses the fetched cloud offset. A cloud 
read with `offsetTso < endTso` can be recorded as current (or a current read as 
old), so `MTMVTask` publishes the wrong snapshot/epoch and can serve stale 
rows. Compute the flag from the finalized cloud state and add MOW/DUP cloud 
coverage.



##########
fe/fe-core/src/main/java/org/apache/doris/catalog/stream/OlapTableStreamWrapper.java:
##########
@@ -363,6 +366,35 @@ public Map<Long, Pair<Long, Long>> 
getHistoryPartitionOffsets(List<Long> selecte
                         s -> Pair.of(null, 
TSOTimestamp.toExclusiveBound(s.getValue().first))));
     }
 
+    /**
+     * Whether a snapshot read of these partitions answers with the table as 
of the stream offset rather than
+     * with the table as it is now.
+     *
+     * <p>It does for the partitions holding data the offset has not consumed: 
those are rebuilt from the
+     * binlog into the image at that offset, while the rest are read directly. 
Which partitions those are
+     * depends on the table's key type, because the read is built differently 
for each -- a merge-on-write
+     * table splits its partitions into the two kinds, and a duplicate-key one 
reads all of them bounded by
+     * the offset. This is the split {@code NormalizeOlapTableStreamScan} 
makes when it binds the read, asked

Review Comment:
   [P1] Do not classify a MOW partition with no consumption baseline as an 
older snapshot. `hasData` returns true when `partitionOffset` is absent and the 
partition has rows, so this `normal` set excludes it; however 
`makeSnapshotScan` first removes no-baseline IDs via 
`filterConsumedPartitionIds`, meaning no historical image was read. The batch 
is nevertheless marked `readFromAStreamOffset` and its snapshots/epochs are 
withheld, potentially leaving a newly populated MV partition empty and forcing 
repeated rebuilds. Base this flag on the IDs actually consumed by the 
normalized snapshot plan and add a no-baseline MOW case.



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