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


##########
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:
   Right. `576743fd404` fixes it, and it takes the key-type branch with it.
   
   `makeSnapshotScan` starts with `filterConsumedPartitionIds`, so a partition 
with no consumption baseline is never read through the stream; the judgement 
asked `filterNormalSnapshotPartitionIds` about the unfiltered ids, and 
`hasData` reads true for such a partition when it holds rows, so it was counted 
as behind and a batch that had read the table as it is withheld its snapshots 
and epochs. The judgement now does what the read does, in the same order: drop 
the unconsumed partitions, then ask which of the rest are read directly.
   
   That filter is also what makes one expression answer for both key types, so 
the `DUP_KEYS` branch is gone. A duplicate-key read is bounded by the offset 
for every partition rather than split into two kinds, and a partition whose 
offset has reached the end of the table is read as the table is either way -- 
the two shapes differ in how the read is built, not in which partitions answer 
with an older image.
   
   `OlapTableStreamSnapshotReadTest` pins the merge-on-write no-baseline case 
now, which is what it was missing: the same test covered it for duplicate keys 
only, and that asymmetry is exactly what let this through. Both key types are 
covered for the behind, caught-up and no-baseline cases.
   
   Verified with the suites around these reads: 
`test_ivm_partial_rebuild_is_not_recorded`, 
`test_ivm_fallback_keeps_the_whole_scope`, 
`test_ivm_fallback_stream_multi_batch` and `..._dup`, and `test_ivm_snapshot` 
as the canary for withholding too much or too little, plus 
`IvmFullRefreshMTMVTest`, `MTMVTaskTest` and the new wrapper test (79 unit 
tests, checkstyle clean).
   



##########
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:
   Right, and it is the point I could not settle from here -- cloud mode is not 
something I can run locally, so thank you for the ordering.
   
   `576743fd404` moves the judgement to where the read is built, 
`NormalizeOlapTableStreamScan`, right after that rule drops the unconsumed 
partitions. That is after `afterAnalyze` has installed the read states (and 
after `IvmIncrRefreshMTMV`'s explicit install for the incremental path, where 
no snapshot scan reaches this anyway), so the wrapper answers from the state 
the read itself will use, and the judgement is the read's own split asked one 
line earlier. The flag still reaches the refresh through the 
`IvmRewriteContext` a refresh owns per statement -- no field on the 
`StatementContext`, no planner hook, and nothing added to the cloud read-state 
path.
   
   "Compute the flag from the finalized cloud state" is what I should have done 
instead of asking the wrapper from the rule that only queues the hook; that 
rule runs before the states are installed, so in cloud mode the wrapper was 
answering from the offsets it holds locally.
   
   Coverage: the two key types are pinned in `OlapTableStreamSnapshotReadTest` 
for behind, caught-up and no-baseline, and the end-to-end suite 
(`test_ivm_partial_rebuild_is_not_recorded`) is what proves the recording still 
fires from where it moved to -- it asserts the rebuilt partition is still 
unrecorded after the move. Cloud itself I cannot exercise here; the states the 
wrapper reads at that point are the ones the read uses, which is what makes the 
answer the same in either mode, and the CI run is the check on that.
   



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