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


##########
fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java:
##########
@@ -792,13 +930,213 @@ private IvmIncrRefreshResult 
executeSingleIvmAttempt(MTMVRefreshContext refreshC
         }
         if (ivmResult.isSuccess()) {
             this.partitionSnapshots.putAll(capturedSnapshots);
-            this.completedPartitions.addAll(needRefreshPartitions);
+            recordRefreshCompleted(incrementalScope);
+            commitCapturedEpochs(capturedEpochs);
             LOG.info("IVM incremental refresh succeeded for mv={}, taskId={}",
                     mtmv.getName(), getTaskId());
         }
         return ivmResult;
     }
 
+    /**
+     * Whether the delta is what a partition this task replaced is still 
waiting for: the records the rebuild
+     * held back, whose publication is the delta's to earn.
+     *
+     * <p>An attempt whose scope comes out empty is one where every partition 
that needs a refresh is one the
+     * rebuild above replaced -- they are taken out of the scope on purpose. 
Skipping the delta then would
+     * skip it for exactly the partitions it is the only repair of: what the 
rebuild read of a table the MV
+     * does not partition by is the image as of the stream offset, and the 
delta is what brings that table up
+     * to date for them. Their records stay unrecorded until it has run, so 
the attempt runs for them even
+     * with nothing to record -- the scope is what it records, and an empty 
one records nothing.
+     *
+     * <p>The delta does not need the scope to do that work: its plan is the 
MV's own query over the streams,
+     * and the partitions it writes are the ones the change reaches, which 
includes the partitions the
+     * rebuild replaced.
+     *
+     * <p>One of the two held maps answers for both: a batch holds its epochs 
and its snapshots together,
+     * and redeeming publishes the pair -- which is also what the two are read 
as.
+     */
+    private boolean hasRecordsHeldForTheDelta() {
+        return !epochsHeldUntilTheDeltaRuns.isEmpty();
+    }
+
+    /**
+     * Adds a phase's scope to what this task reports as refreshed, and the 
partitions it committed to what
+     * this task reports as done. Both are the task's; see the fields for why.
+     */
+    private void recordRefreshScope(Collection<String> partitions) {
+        // Created by the first phase that records one: a task that has not 
refreshed anything reports
+        // nothing, which is the same state as one that has not run yet. 
Concurrent because the columns it
+        // feeds are read while the worker fills them -- the tasks() table 
function reports a running task
+        // -- and ordered because what this reports is persisted and has to 
read the same on every look.
+        if (needRefreshPartitions == null) {
+            needRefreshPartitions = new ConcurrentSkipListSet<>();
+        }
+        needRefreshPartitions.addAll(partitions);
+    }
+
+    private void recordRefreshCompleted(Collection<String> partitions) {
+        if (completedPartitions == null) {
+            completedPartitions = new ConcurrentSkipListSet<>();
+        }
+        completedPartitions.addAll(partitions);
+    }
+
+    /**
+     * Says durably that the parts of the scope that do not name a rebuild 
requirement have to be rebuilt, and
+     * brings the epochs this phase records in line with it.
+     *
+     * <p>Raised here rather than by each caller of the executor, like the 
scope above: this is the phase that
+     * replaces partitions, and a caller that forgot would leave a partition 
whose rows were never published
+     * looking caught up. The partitions a caller has already made dirty are 
left as they are, which is what
+     * makes raising it here harmless for them: a whole-MV attempt marks its 
scope before it reconciles the
+     * streams, and the incremental attempt rebuilds the partitions an 
invalidation marked.
+     * See MTMV#raiseRebuildRequirement.
+     *
+     * <p>What that call reports is what a partition it raised now names, and 
this phase is clamped to it. The
+     * clamp cannot stay at what an earlier attempt planned: that value sits 
below the requirement this phase
+     * has just raised, so the epochs recorded here would leave the partition 
dirty after it was replaced, and
+     * every refresh after it would rebuild the same partitions again.
+     *
+     * <p>A partition that already named a requirement keeps the entry it has, 
and is deliberately not moved
+     * up to what it names now. The entry is the value the routing decision 
saw, and it is what keeps a mark
+     * landing between that decision and this phase's read from being recorded 
as met by a replacement that
+     * read before the change it made. Leaving it where it is costs one 
rebuild; moving it up could cost the
+     * change.
+     */
+    private void raiseRequirementForRefreshScope(Collection<String> 
partitions) {
+        if (!mtmv.isIvm()) {
+            // A plain MV has no streams to read and its epochs record 
nothing: the sync criterion plans it
+            // again on its own until its snapshots are published.
+            return;
+        }
+        
ivmPlannedEpochs.putAll(mtmv.raiseRebuildRequirement(Sets.newHashSet(partitions)));
+    }
+
+    /**
+     * Brings the routing decision up to date after a retry has synchronized 
and aligned the MV's partitions.
+     *
+     * <p>A partition the alignment creates is dirty by construction -- {@code 
{0, 1}}, behind its
+     * requirement -- and the decision was taken before it existed. Reading 
the states again is what makes
+     * the retried attempt treat it as such: it joins the dirty set, so the 
incremental attempt leaves it
+     * out rather than recording what a delta captured as the partition being 
caught up, and it gets the
+     * entry the batches are clamped against, so a mark landing later in this 
task cannot be written back as
+     * satisfied either.
+     *
+     * <p>A partition this leaves dirty is not rebuilt here. The rebuild phase 
has already run, and the
+     * partition the alignment created holds no rows yet -- there is nothing 
to replace in it -- so what it
+     * needs is a build, which is what the next refresh's rebuild gives it.
+     *
+     * <p>The planned value of a partition that already had one is kept: that 
is the value the routing
+     * decision was made on, which is what the clamp is for.
+     */
+    private void adoptPartitionsCreatedByTheRetry(Set<String> dirtyPartitions) 
{
+        Set<String> livePartitionNames = mtmv.getPartitionNames();
+        for (Entry<String, MTMVPartitionState> entry : 
mtmv.getPartitionStates().entrySet()) {
+            if (!livePartitionNames.contains(entry.getKey())) {
+                continue;
+            }
+            ivmPlannedEpochs.putIfAbsent(entry.getKey(), 
entry.getValue().getLatestEpoch());
+            if (entry.getValue().isDirty()) {
+                dirtyPartitions.add(entry.getKey());
+            }
+            // What this task holds for a name is that name's only while the 
partition behind it is the same
+            // one: the retry's partition sync can drop a partition and add 
another of the same name back, and
+            // then the captures and snapshots describe a partition that is 
gone. Writing them back would
+            // credit the new one with what the old one held, which is worse 
than a wrong number -- a partition
+            // clean at an epoch a later change only raises to is one no 
refresh rebuilds, so the rows the
+            // recreation removed would be published as current.
+            //
+            // Told apart by id rather than by the state: a partition this 
task rebuilt has its rows and its
+            // capture, and its state still reads as the one an entry starts 
with until the task result writes
+            // the epochs back, so the state cannot say whether the name means 
the same partition.
+            // A partition the sync dropped has no state to walk here, so this 
one is live; it could still be
+            // dropped by a concurrent DDL, which is the one case where there 
is nothing to compare with.
+            Partition partition = mtmv.getPartition(entry.getKey());
+            Long capturedId = capturedPartitionIds.get(entry.getKey());

Review Comment:
   Right, and the reason the guard did not fire is in which path records the 
id: `commitCapturedEpochs` is what writes `capturedPartitionIds`, and a batch 
whose records are held back does not go through it -- it puts them in the two 
held maps and stops. So `capturedId` was null here for exactly the names the 
fence has to be able to drop, and the fence read that as "nothing to compare 
with".
   
   `cafc968c0c09` fixes both halves. Holding a batch's records now records the 
partition each was taken under, through the same `recordCapturedPartitionId` 
the committing path uses, and the fence drops the held records along with the 
committed ones -- both maps, since what a later delta would redeem is the pair.
   
   One thing I corrected while reviewing it: the ids are recorded for the union 
of the two maps' names, not the epochs alone. `captureLatestEpochs` answers 
only for partitions that already had an entry, while the snapshots are 
generated for the whole batch, so a partition can be held with a snapshot and 
no epoch -- and naming the epochs would have left that one without an id, which 
is the same hole one case narrower.
   
   `MTMVTaskTest` pins both halves and both directions: a held batch recording 
the ids it was taken under (including the snapshot-without-an-epoch name) while 
publishing nothing, and the fence dropping the held records of a name whose 
partition changed while keeping those of a name whose partition did not. Each 
fails with its half reverted -- the first on the recorded ids, the second 
leaving the held entries in place.
   
   On the two-retry schedule: I have not built it end to end, for the reason 
the earlier retry thread gives -- it needs an `MV_PARTITION_NOT_FOUND` whose 
sync drops and re-adds an MV partition *inside* the retry window, and I have 
not found a way to place a DDL in that window from a suite. What makes the 
fence complete without it is where replacements can come from: tasks serialize 
under the MTMVJob lock, the sync at the start of a refresh runs before any 
capture, and an IVM MV's partitions are not user-DDL'd -- so the retry's sync 
is the only sync that can replace a partition after this task captured 
anything. What remains uncovered, and stays documented at the guard, is a 
capture taken while no partition of that name was live: there is no id to 
compare with, and that one predates this commit.
   



##########
fe/fe-core/src/main/java/org/apache/doris/catalog/stream/OlapTableStreamWrapper.java:
##########
@@ -363,6 +365,41 @@ public Map<Long, Pair<Long, Long>> 
getHistoryPartitionOffsets(List<Long> selecte
                         s -> Pair.of(null, 
TSOTimestamp.toExclusiveBound(s.getValue().first))));
     }
 
+    /**
+     * Whether the snapshot read of these partitions answers with the table as 
it is now.
+     *
+     * <p>It does not when the read leaves out rows the table holds, and that 
happens in two ways. The read
+     * drops the partitions with no consumption baseline -- no offset, or the 
sentinel one of a partition that
+     * was empty when the stream was created -- because there is no offset to 
read them from; rows such a
+     * partition holds by now are rows the answer does not have. And the 
partitions it does read it reads as
+     * the table is only where the offset reached the end of the partition: 
what is behind that offset it
+     * answers with the image at the offset, which is the table as it was 
then. This is the question
+     * {@code NormalizeOlapTableStreamScan} answers when it binds the read, 
asked by the refresh that may not
+     * record a partition it answered from an incomplete or older image as 
holding the table's current state.
+     *
+     * <p>Asked in one expression for both key types: a duplicate-key read is 
bounded by the offset for every
+     * partition rather than split into the two kinds, and a partition whose 
offset reached the end is read as
+     * the table is by either of them.
+     */
+    public boolean answersWithTheCurrentTable(List<Long> partitionIds) {
+        List<Long> consumed = filterConsumedPartitionIds(partitionIds);
+        Set<Long> consumedIds = ImmutableSet.copyOf(consumed);
+        Set<Long> atTheEnd = 
ImmutableSet.copyOf(filterNormalSnapshotPartitionIds(consumed));
+        for (Long partitionId : partitionIds) {
+            if (!consumedIds.contains(partitionId)) {
+                // Not read at all: rows it holds are missing from the answer.
+                if (getBaseTable().getPartition(partitionId).hasData()) {

Review Comment:
   This one I would rather not change, because `CloudPartition.hasData()` does 
not answer from the cache in the case you describe. The cache is a 
short-circuit for the positive: `hasDataCached()` returning true answers true, 
and anything else -- cached-empty, unknown, cache disabled -- falls through to 
`getVisibleVersion()`, which issues the `get_version` RPC to the meta service 
and writes the answer back into the cache 
(`cloud/catalog/CloudPartition.java:519-533`, and the RPC path it lands in at 
`:175-210`). A partition whose rows have committed is answered from the meta 
service's version, not from a stale cached 1.
   
   The direction the cache can skew this is the opposite one: a stale *hit* 
(`visible version > 1` from a cache that a drop has not caught up with yet) 
makes `hasData()` answer true, which withholds more than it needs to rather 
than less. And with `disable_empty_partition_prune` set, `hasDataCached()` is 
true for everything, so nothing is dropped for being empty at all.
   
   So for `answersWithTheCurrentTable` to return true while the read dropped 
populated rows, `hasData()` would have to return false with rows present -- 
which needs the meta service to answer `VERSION_NOT_FOUND` or version 1 for a 
partition whose rows are committed. If you have a path to that, the line to 
point at is the comparison in `hasData()`; I have not found one, and the read 
state's own `visible_version` would be a staler source than the version this 
asks for, so I would not swap the source either.
   
   What the branch does cost is an RPC per dropped partition at planning time, 
which is bounded by the partitions that have no consumption baseline at all -- 
the case this judgement is about. Tables that have consumed their partitions 
reach the cached-true or cached-false-with-data path without one.
   



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