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


##########
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:
   [P1] Classify omitted cloud partitions from the installed read state. 
`filterConsumedPartitionIds` can omit an initially empty non-PCT dimension 
partition because its stream offset has no real baseline, while pruning kept it 
using the installed state's visible version > 1. After rows commit, this 
`CloudPartition.hasData()` call can still return cached version 1 before 
asynchronous cache invalidation reaches this FE. The SNAPSHOT plan then drops 
those rows while `answersWithTheCurrentTable` returns true, so a partial 
overwrite can publish a clean epoch/current snapshot even if the following 
strict delta fails. Use the installed state's visible version for this cloud 
branch and test installed version > 1 with cached version 1. The existing 
no-baseline and hook-timing threads do not cover this version-source mismatch.



##########
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:
   [P1] Clear held records when retry synchronization replaces an MV partition. 
A partial non-PCT SNAPSHOT rebuild can put p's epoch and snapshot only in the 
two `*HeldUntilTheDeltaRuns` maps, so `capturedPartitionIds[p]` is null here. 
If one retry sync drops p and a later sync re-adds a new p under the same name, 
this guard leaves both held entries intact. A succeeding delta then calls 
`redeemHeldRecords`, and `commitCapturedEpochs` binds the old epoch to the new 
ID. `ADD_TASK` can turn new p from `{0,1}` into clean `{2,1}` and publish the 
old snapshot even though that delta never built its baseline. Record the 
partition ID for held work and discard both held entries on replacement; cover 
the two-retry, successful-delta schedule. The existing retry thread covers 
already committed captures, whereas these records have no captured ID yet.



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