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]