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]