yujun777 commented on code in PR #68390:
URL: https://github.com/apache/doris/pull/68390#discussion_r4122177634
##########
fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java:
##########
@@ -928,8 +1217,26 @@ private void
executePartitionBasedRefresh(MTMVRefreshContext context, RefreshMod
mtmv.getName(), getTaskId(), e);
throw new JobException(e.getMessage(), e);
}
- completedPartitions.addAll(execPartitionNames);
- partitionSnapshots.putAll(execPartitionSnapshots);
+ recordRefreshCompleted(execPartitionNames);
+ // What this batch may record: everything it replaced, unless the
read behind it answered with a
+ // base table the MV does not partition by as of an older state
than that table is in now -- which
+ // a partial read does for the partitions holding data the stream
offset has not consumed, and it
+ // is the delta that follows that brings the table up to date. For
the partitions this batch
+ // replaced, that delta does not apply: they are kept out of its
scope, so a delta computed against
+ // the rows a rebuild replaced cannot be applied twice. Recording
them would say they hold the
+ // table's current state when they hold that image, and nothing
would plan them again. Left
+ // unrecorded, the requirement stays, the MV keeps the snapshot it
has rather than one claiming the
+ // current state, and a later refresh rebuilds them -- by then the
offset has been consumed, so a
+ // rebuild reads the table as it is. What the read answered with
is recorded where the read is
+ // planned; see IvmFullRefreshMTMV#readsAnOlderImage.
+ if
(rewriteContext.map(IvmRewriteContext::isReadFromAStreamOffset).orElse(false)) {
Review Comment:
Right, and the evidence is in the suite my last commit added -- I read `kept
out of its scope` as `the delta does not write to it`, when the scope only
decides what the task records. The delta's target is the partitions the change
touches, and that includes the ones the rebuild replaced; my own case showed
the rebuilt partition moving from `10` to `20` during the successful delta
while I was still withholding its epoch and snapshot. Holding them was
therefore a full rebuild of those partitions on every refresh, for as long as a
table they read keeps changing -- which is what you describe.
`635c078bcb4` holds them instead of dropping them, and publishes them when
the incremental attempt's delta succeeds. Two things I checked before doing
that, since it is the direction that can publish too much:
* the reason a batch is held is always a table the MV does not partition by
-- the flag is set only where a snapshot read is bound, and the table the MV
partitions by is read with `RESET` -- so the delta that redeems it is that
table's delta, over the partitions the change maps to;
* a held partition the change does not touch is one whose rows the change
does not reach, so publishing it is not claiming rows it does not have.
When the delta does not run or fails, the records stay held, which is the
case the holding exists for. The visible effect in the suite added with the
last commit: its second pass now reports no scope and no rebuild where it used
to report one partition and one, so the recovery is finished there rather than
deferred by one -- the golden is regenerated and reviewed line by line for that.
##########
fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java:
##########
@@ -792,13 +912,159 @@ 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;
}
+ /**
+ * 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());
+ }
+ // A partition the alignment has just created -- the state an
entry starts with,
+ // MTMVPartitionState.initial() -- is a
+ // partition of this name that the retry's partition sync
recreated: the partition the captures and
+ // snapshots this task holds for that name describe is gone, and
the one that took its name holds
+ // nothing. Writing those 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.
+ if (entry.getValue().getRefreshEpoch() == 0 &&
entry.getValue().getLatestEpoch() == 1) {
Review Comment:
Right, and the shape of the state is exactly what cannot tell the two apart:
a partition this task rebuilt has its rows and its capture committed, and its
state still reads the one an entry starts with until the task result writes the
epochs back -- so my fence took the partition I had just repaired for one that
had been recreated, and dropped what it had committed. Leaving a repaired
partition unrecorded is not free either: it is rebuilt in full by the next
refresh, and out of transparent rewrite until then.
`635c078bcb4` fences by identity instead: `commitCapturedEpochs` records the
id of the MV partition the capture was taken under, and
`adoptPartitionsCreatedByTheRetry` drops the capture, snapshot and record only
when the live partition of that name carries a different id -- which is what a
sync drop-and-re-add leaves behind. A partition the task itself rebuilt keeps
its id, so what it committed stands.
`MTMVTaskTest` covers both directions now: same id with a
freshly-aligned-looking state keeps everything, different id drops all three.
The end-to-end case you describe -- a newly aligned `p`, a `q` arriving later,
the retry syncing `q` and leaving `p` alone -- I have not built; it needs
repeatedly inducing a `MV_PARTITION_NOT_FOUND` delta failure by dropping and
re-adding populated base partitions, which is why I pinned the mechanism at the
unit level. Say the word and I will add it as its own suite rather than growing
one that already drives several shapes.
--
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]