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]