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]