yujun777 commented on code in PR #68390:
URL: https://github.com/apache/doris/pull/68390#discussion_r4122827607
##########
fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java:
##########
@@ -715,13 +796,52 @@ private AttemptResultType
executeIvmAttempt(MTMVRefreshContext refreshContext,
+ "Continuing with COMPLETE refresh.", mtmv.getName(),
getTaskId());
return AttemptResultType.FALLBACK_TO_COMPLETE;
}
+ // The partitions the criterion says must be rebuilt rather than
caught up: the delta path can only
+ // append, so a partition it treated as current would record that in
its epoch while its rows still
+ // come from before the change. Rebuilt first, with the partition
executor, because that is the
+ // full recomputation they need -- and only in this task's batches, so
a change that arrives while
+ // it runs leaves them dirty for the next round instead of being
swallowed.
+ // One read of the states decides both what has to be rebuilt and the
requirement each batch may
+ // write back. Reading them separately would leave a window between
the two in which a mark lands,
+ // the routing decision does not see it, and the batch that follows
captures the raised requirement
+ // and records it as met by a delta that cannot remove the rows that
mark made unusable.
+ Map<String, MTMVPartitionState> plannedStates =
mtmv.getPartitionStates();
+ Set<String> livePartitionNames = mtmv.getPartitionNames();
+ Set<String> dirtyPartitions = Sets.newLinkedHashSet();
+ Map<String, Long> plannedEpochs = Maps.newHashMap();
+ for (Entry<String, MTMVPartitionState> plannedState :
plannedStates.entrySet()) {
+ if (!livePartitionNames.contains(plannedState.getKey())) {
+ continue;
+ }
+ plannedEpochs.put(plannedState.getKey(),
plannedState.getValue().getLatestEpoch());
+ if (plannedState.getValue().isDirty()) {
+ dirtyPartitions.add(plannedState.getKey());
+ }
+ }
+ this.ivmPlannedEpochs = plannedEpochs;
+ if (!dirtyPartitions.isEmpty()) {
+ LOG.info("Rebuilding {} invalidated MV partitions before the
incremental refresh, mv={}, taskId={}",
+ dirtyPartitions.size(), mtmv.getName(), getTaskId());
+ List<String> toRebuild = Lists.newArrayList(dirtyPartitions);
+ toRebuild.sort(Comparator.naturalOrder());
+ this.refreshMode = generateRefreshMode(toRebuild);
+ try {
+ executePartitionBasedRefresh(refreshContext,
RefreshMode.PARTITIONS, ctx, toRebuild);
+ } finally {
+ recordRebuiltPartitions(request);
+ }
+ }
MTMVRefreshContext currentRefreshContext = refreshContext;
int ivmAttemptLimit = Math.max(Config.max_query_retry_time, 0) + 1;
IvmIncrRefreshResult ivmResult = null;
for (int partitionSyncRetryCount = 0;
partitionSyncRetryCount < ivmAttemptLimit;
partitionSyncRetryCount++) {
- ivmResult = executeSingleIvmAttempt(currentRefreshContext);
+ ivmResult = executeSingleIvmAttempt(currentRefreshContext,
dirtyPartitions);
if (ivmResult.isSuccess()) {
+ // The delta has run, and it is what brought the tables those
partitions were rebuilt from up
+ // to date: its target is the partitions the change touches,
which includes the ones the
+ // rebuild replaced. What the rebuild wrote is current now, so
what was held can be recorded.
+ redeemHeldRecords();
Review Comment:
Right, and the step it skips is exactly what the publication is for.
`redeemHeldRecords` runs on `ivmResult.isSuccess()`, and
`executeSingleIvmAttempt` returns success without reaching `doRefresh` when its
scope comes out empty -- and the scope is *what the attempt records*, emptied
by taking the partitions the rebuild replaced out of it, not a statement that
there is nothing left to consume. It does not even need the retry to happen:
the retry is simply the shape the code's own test already pins as reachable,
since `adoptPartitionsCreatedByTheRetry` adds the recreated partition to the
dirty set, so the scope can empty on the second attempt while the first
attempt's delta failed.
`42ed8c44575` runs the delta for the records it holds back: the skip now
needs the scope to be empty *and* nothing held. `hasRecordsHeldForTheDelta`
states why the two questions differ, and why the delta does not need the scope
to consume the change -- 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. The attempt still records nothing (the scope
is what it records), and what was held is published by the delta's success.
The unit test pins both directions:
`MTMVTaskTest#testTheIncrementalAttemptRunsForTheRecordsAPartialReadHeldBack`
drives the state you describe -- held records, every refresh-needing partition
dirty, so the scope is empty -- and asserts that `doRefresh` ran and that the
held records were then published; `...IsSkippedWhenThereIsNothingHeldBack` pins
that an attempt with nothing held is still skipped.
The end-to-end retry-to-empty sequence is still not built: it needs a
`MV_PARTITION_NOT_FOUND` that the following sync resolves, and I have not found
a recipe that induces it reliably rather than by timing. Say the word and I
will add it as its own suite rather than pinning the mechanism at the unit
level only.
##########
fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java:
##########
@@ -444,7 +499,33 @@ private RefreshRequest resolveRefreshRequest() throws
JobException {
Lists.newArrayList(), false);
}
- private List<RefreshAttemptType> buildAttempts(RefreshRequest request,
boolean containsOneRowRelation) {
+ private List<RefreshAttemptType> buildAttempts(RefreshRequest request,
boolean containsOneRowRelation)
+ throws JobException {
+ // A schema-level invalidation is not a set of dirty partitions: it
means every partition, including
+ // the ones partition sync has not created yet, and no per-partition
requirement can express that.
+ // IVM only -- a non-IVM MV reaches the same effect through its
cleared snapshot, which its own
+ // refresh already depends on.
+ //
+ // Judged before the initial-refresh shortcut below, which also
answers COMPLETE: an MV that has
+ // never been refreshed and reads an excluded trigger table (or a
one-row relation) has to be built
+ // by a whole-MV refresh, and a request that may not fall back has to
hear that rather than have it
+ // decided for it -- otherwise the refusal here would be unreachable
in exactly the state it names.
+ if (mtmv.isIvm() && !request.explicitPartitions
+ && mtmv.getStatus().getState() == MTMVState.SCHEMA_CHANGE) {
+ if (request.refreshMode == RefreshMode.PARTITIONS &&
!request.allowFallback) {
Review Comment:
This one is by design, and the widening is reported rather than silent: the
task's `RefreshMode` reads COMPLETE after a strict `REFRESH ... INCREMENTAL`
met an invalidated MV. SCHEMA_CHANGE is a property of the MV, not a scope the
request named -- its own baseline is gone, and its query may no longer analyze
at all (a rename is what the suite uses) -- so nothing incremental can repair
it, and the COMPLETE refresh restores exactly the baseline the incremental
attempt would have read. Refusing would leave a strict caller in a state that
only a COMPLETE can leave, which is what the message on the PARTITIONS branch
already tells them to run.
That is a different question from the one the all-dirty and unusable-stream
shortcuts answer. Those guard a strict request against a widening it did not
ask for -- rebuilding partitions it excluded, or a stream reset its incremental
attempt would have failed on instead. Here the state, not the request's scope,
is what forces the rebuild, and the result says so.
`MTMVTaskTest#testBuildAttemptsRebuildsTheWholeMvForAStrictIncrementalInSchemaChange`
and case 3 of `test_ivm_partition_epoch_rebuild` pin it, including the rebuilt
count a request that did not ask for the rebuild reports.
##########
fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java:
##########
@@ -738,28 +1067,61 @@ public boolean invalidateIvmBaseline(BaseTableInfo
baseTableInfo, Map<String, Lo
+ "changedPartitions={}", name, baseTableInfo,
changedPartitions);
return false;
}
+ if (!affectedMvPartitions.isPresent()) {
+ // A narrower rebuild could leave a partition holding rows of the
changed base partition
+ // untouched, and those rows cannot be repaired later: the change
emitted no row binlog. The
+ // whole MV is invalidated instead, which says "every partition,
including the ones partition
+ // sync has not created yet" -- what a per-partition requirement
cannot express.
+ invalidateWholeMv(reason).await();
+ return true;
+ }
EditLogItem editLogItem;
writeMvLock();
try {
- if (ivmInfo == null) {
- ivmInfo = new IvmInfo();
- }
- if (!affectedMvPartitions.isPresent()) {
- // A narrower rebuild could leave a partition holding rows of
the changed base partition
- // untouched, and those rows cannot be repaired later: the
change emitted no row binlog.
- ivmInfo.requireCompleteBaselineRebuild();
- } else {
-
ivmInfo.addPendingBaselineRebuildPartitions(affectedMvPartitions.get());
+ // Placed: the partitions that read the change get the requirement
raised, which is what sends
+ // them to a rebuild while every other partition keeps catching up
incrementally. No version
+ // bump here -- a partial invalidation does not invalidate a task
result, and the requirement it
+ // raises survives the write-back by construction.
+ Set<String> marked =
markIvmPartitionsInvalidated(affectedMvPartitions.get());
+ if (marked.isEmpty()) {
+ LOG.debug("No MV partition holds the changed base partitions,
mv={}, baseTable={}, "
+ + "changedPartitions={}", name, baseTableInfo,
changedPartitions);
+ return false;
}
- schemaChangeVersion++;
- editLogItem = submitIvmInfoChange();
+ editLogItem = submitPartitionStatesChange(marked);
Review Comment:
Right, and the shape is worse than the size: this is the base table's DDL
path, under the MV write lock, so one `TRUNCATE` of a single base partition
copies and journals every entry of an MV whose states the change did not touch.
The map was also the caller's own by reference, so what the record carried was
the live field.
`42ed8c44575` writes the marked entries as the delta they are -- detached
copies, the shape `raiseRebuildRequirement` already uses -- with the snapshot
removals riding in the same record, which the replay merges over the states
other records wrote. Alignment and the whole-MV mark keep the map, since those
are the changes that do move every entry.
`IvmBaselineRebuildTest#testAPartitionScopedInvalidationJournalsOnlyWhatItMarked`
captures the payload of a real `DROP PARTITION` (the edit log spied,
delegating to the real one), asserts it carries exactly the marked partition
with `merge` set and the removal, and then replays it over a state map that
another record has moved in between -- so a payload that replaced the map
instead of merging fails.
##########
fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java:
##########
@@ -663,54 +746,299 @@ public void alterIvmInfo(IvmInfo ivmInfo) {
* into the journal, and a replay that replaces the field would leave the
caller's reference
* pointing at state that is no longer the MV's. Changing the states is
the MV's own job, under its
* write lock.
- *
- * <p>A missing map -- an image written before the field existed, or a
non-IVM MV -- reads as empty.
*/
public Map<String, MTMVPartitionState> getPartitionStates() {
readMvLock();
try {
- if (partitionStates == null) {
- return Collections.emptyMap();
- }
return
Collections.unmodifiableMap(MTMVPartitionState.copyOf(partitionStates));
} finally {
readMvUnlock();
}
}
+ /**
+ * The partitions whose requirement has been raised and not met, which a
refresh has to rebuild rather
+ * than catch up.
+ *
+ * <p>Detached names rather than the states themselves: a caller that only
routes by them has no
+ * business holding the map the MV journals, and it needs nothing else
from an entry.
+ */
+ public Set<String> getPartitionsNeedingRebuild() {
+ // Built before the lock, like the map getLatestEpochs returns: which
entries go in is what needs
+ // the lock, not having somewhere to put them.
+ Set<String> res = Sets.newLinkedHashSet();
+ readMvLock();
+ try {
+ for (Entry<String, MTMVPartitionState> entry :
partitionStates.entrySet()) {
+ if (entry.getValue().isDirty()) {
+ res.add(entry.getKey());
+ }
+ }
+ return res;
+ } finally {
+ readMvUnlock();
+ }
+ }
+
+ /**
+ * Whether every partition the MV holds needs a rebuild, which is when a
whole-MV refresh does nothing
+ * the per-partition routing would not.
+ *
+ * <p>A partition that holds data and does not need one makes this false:
a whole-MV refresh would
+ * recompute it for nothing, which is the waste the per-partition routing
exists to avoid. A partition
+ * that was never refreshed does not count against it -- a whole-MV
refresh fills it, which its
+ * per-partition branch would do as well -- and it needs no clause of its
own: an aligned entry is
+ * {@code {0, 1}}, so it is behind its requirement already. An MV with no
partitions is not an
+ * escalation either.
+ *
+ * <p>Read in place rather than through {@link #getPartitionStates()}: the
caller asks a yes/no
+ * question, and copying the map to answer it would allocate a state
object per partition, under this
+ * lock, on every refresh -- including the ones that escalate nothing.
+ */
+ public boolean allPartitionsNeedRebuild() {
+ readMvLock();
+ try {
+ return !partitionStates.isEmpty()
+ &&
partitionStates.values().stream().allMatch(MTMVPartitionState::isDirty);
+ } finally {
+ readMvUnlock();
+ }
+ }
+
// ALTER_PARTITION_STATES replay applies a detached snapshot here,
mirroring alterIvmInfo(). Live
// invalidation changes submit their journal from the mutating method
instead.
//
// A payload without the member carries no state at all, which is not the
same as an empty map that
// says the states are now empty: leaving them alone is the only answer
that cannot lose state.
public void alterPartitionStates(Map<String, MTMVPartitionState>
partitionStates) {
- if (partitionStates == null) {
+ replayAlterPartitionStates(partitionStates, null, false);
+ }
+
+ /**
+ * ALTER_PARTITION_STATES replay: applies the states the payload carries,
and drops the snapshots it
+ * names. Both in one lock acquisition, because a reader that saw the new
requirement while the
+ * snapshot was still there could let a transparent rewrite serve rows the
rebuild has to replace.
+ *
+ * <p>A payload without the states carries none, which is not the same as
an empty map that says the
+ * states are now empty: leaving them alone is the only answer that cannot
lose state.
+ *
+ * <p>{@code merge} says what the payload's states are. False, which is
what a payload written before the
+ * member existed means and what the invalidation channel still writes,
carries the map itself and
+ * replaces. True carries only the partitions a change touched -- see
submitPartitionStatesDelta -- so the
+ * entries it does not name belong to other records (an invalidation that
ran during the refresh, an entry
+ * alignment added) and are merged over rather than dropped.
+ */
+ public void replayAlterPartitionStates(Map<String, MTMVPartitionState>
partitionStates,
+ Set<String> removedSnapshotPartitions, boolean merge) {
+ writeMvLock();
+ try {
+ if (partitionStates != null) {
+ if (merge) {
+ for (Entry<String, MTMVPartitionState> entry :
partitionStates.entrySet()) {
+ this.partitionStates.put(entry.getKey(), new
MTMVPartitionState(entry.getValue()));
+ }
+ } else {
+ this.partitionStates =
MTMVPartitionState.copyOf(partitionStates);
+ }
+ }
+ refreshSnapshot.removeSnapshots(removedSnapshotPartitions);
+ } finally {
+ writeMvUnlock();
+ }
+ }
+
+ /**
+ * The {@code latestEpoch} of the given MV partitions, taken under the MV
read lock.
+ *
+ * <p>This is the value a refresh has to remember: what it read from the
base tables is described by
+ * the requirement in force when it started reading, so writing that value
back as the new
+ * {@code refreshEpoch} is what keeps an invalidation arriving mid-refresh
from being swallowed. A
+ * partition without an entry is left out -- a caller writes an epoch only
for what it captured.
+ */
+ public Map<String, Long> getLatestEpochs(Set<String> partitionNames) {
+ if (CollectionUtils.isEmpty(partitionNames)) {
+ return Collections.emptyMap();
+ }
+ // Sized before the lock: the state map is what needs it, and building
the map is not part of that.
+ Map<String, Long> res =
Maps.newHashMapWithExpectedSize(partitionNames.size());
+ readMvLock();
+ try {
+ for (String partitionName : partitionNames) {
+ MTMVPartitionState state = partitionStates.get(partitionName);
+ if (state != null) {
+ res.put(partitionName, state.getLatestEpoch());
+ }
+ }
+ return res;
+ } finally {
+ readMvUnlock();
+ }
+ }
+
+ /**
+ * Brings the partition states in line with the MV's partitions: every
partition gets an entry, and
+ * every entry whose partition is gone is dropped.
+ *
+ * <p>Alignment is what makes "the partition exists" and "the entry
exists" the same thing, and it is
+ * why an invalidation cannot miss: rows are only written by a refresh,
and every refresh aligns
+ * before it reads a base table, so a partition that holds rows always has
an entry for the mark to
+ * land on. The other direction is what makes the criterion safe -- an
entry created here describes a
+ * partition with no rows yet, so requiring one generation of it discards
no requirement that was
+ * made earlier.
+ *
+ * <p>What it changes is journaled, because the entry has to be on disk
before the rows it describes
+ * can be: a crash between this and the task result would otherwise leave
a partition that holds rows
+ * with no entry at all, and every later invalidation of it would find
nothing to land on. That is the
+ * one shape in which the criterion cannot be read -- "no entry" is
supposed to mean "no rows" -- so
+ * the entry is made durable before any base table is read rather than
derived again on the next run.
+ *
+ * <p>It is deliberately not a hook on every path that creates or drops a
partition. An entry is
+ * derived state, and rebuilding it from the live partition set also
repairs whatever a crash left
+ * behind: the drop of a partition and the removal of its entry are two
journal records, and only
+ * their order -- partition first -- is safe, which leaves at most a stale
entry that the next
+ * alignment drops.
+ *
+ * <p>Only an IVM MV is aligned. For a non-IVM MV the map stays as it is,
and every reader treats
+ * "empty" and "no state" the same.
+ */
+ public void alignPartitionStates() {
+ if (!isIvm()) {
return;
}
+ EditLogItem editLogItem = null;
writeMvLock();
try {
- this.partitionStates = MTMVPartitionState.copyOf(partitionStates);
+ // Read here rather than handed in by the caller: a caller has to
read the names before it takes
+ // this lock, and a partition created in between -- by a
concurrent refresh's partition sync --
+ // would then be dropped by the retainAll below, taking with it
the state a following
+ // invalidation has to land on. The read is cheap and takes no
lock of its own, so doing it here
+ // does not add an edge to the lock order.
+ Set<String> livePartitions = Sets.newHashSet(getPartitionNames());
+ boolean changed =
partitionStates.keySet().retainAll(livePartitions);
+ for (String partitionName : livePartitions) {
+ if (!partitionStates.containsKey(partitionName)) {
+ partitionStates.put(partitionName,
MTMVPartitionState.initial());
+ changed = true;
+ }
+ }
+ if (changed) {
+ editLogItem =
submitPartitionStatesChange(Collections.emptySet());
Review Comment:
Leaving this one as it stands. Alignment reads the live partition names and
reconciles the map as a whole (`retainAll` plus a fill), so the record it
writes is the map that reconciliation produced -- and the removal half has to
be expressible, which a delta cannot say: an entry a merge payload omits is a
no-op, not a drop. Splitting the two halves to make the add-only case a delta
puts a branch in the reconciliation for an event that only fires when a
partition is created or dropped, while the invalidation path you flagged above
is the one that ran on every base partition DDL -- which is why that one is
fixed.
##########
fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java:
##########
@@ -663,54 +746,299 @@ public void alterIvmInfo(IvmInfo ivmInfo) {
* into the journal, and a replay that replaces the field would leave the
caller's reference
* pointing at state that is no longer the MV's. Changing the states is
the MV's own job, under its
* write lock.
- *
- * <p>A missing map -- an image written before the field existed, or a
non-IVM MV -- reads as empty.
*/
public Map<String, MTMVPartitionState> getPartitionStates() {
readMvLock();
try {
- if (partitionStates == null) {
- return Collections.emptyMap();
- }
return
Collections.unmodifiableMap(MTMVPartitionState.copyOf(partitionStates));
} finally {
readMvUnlock();
}
}
+ /**
+ * The partitions whose requirement has been raised and not met, which a
refresh has to rebuild rather
+ * than catch up.
+ *
+ * <p>Detached names rather than the states themselves: a caller that only
routes by them has no
+ * business holding the map the MV journals, and it needs nothing else
from an entry.
+ */
+ public Set<String> getPartitionsNeedingRebuild() {
+ // Built before the lock, like the map getLatestEpochs returns: which
entries go in is what needs
+ // the lock, not having somewhere to put them.
+ Set<String> res = Sets.newLinkedHashSet();
+ readMvLock();
+ try {
+ for (Entry<String, MTMVPartitionState> entry :
partitionStates.entrySet()) {
+ if (entry.getValue().isDirty()) {
+ res.add(entry.getKey());
+ }
+ }
+ return res;
+ } finally {
+ readMvUnlock();
+ }
+ }
+
+ /**
+ * Whether every partition the MV holds needs a rebuild, which is when a
whole-MV refresh does nothing
+ * the per-partition routing would not.
+ *
+ * <p>A partition that holds data and does not need one makes this false:
a whole-MV refresh would
+ * recompute it for nothing, which is the waste the per-partition routing
exists to avoid. A partition
+ * that was never refreshed does not count against it -- a whole-MV
refresh fills it, which its
+ * per-partition branch would do as well -- and it needs no clause of its
own: an aligned entry is
+ * {@code {0, 1}}, so it is behind its requirement already. An MV with no
partitions is not an
+ * escalation either.
+ *
+ * <p>Read in place rather than through {@link #getPartitionStates()}: the
caller asks a yes/no
+ * question, and copying the map to answer it would allocate a state
object per partition, under this
+ * lock, on every refresh -- including the ones that escalate nothing.
+ */
+ public boolean allPartitionsNeedRebuild() {
+ readMvLock();
+ try {
+ return !partitionStates.isEmpty()
+ &&
partitionStates.values().stream().allMatch(MTMVPartitionState::isDirty);
+ } finally {
+ readMvUnlock();
+ }
+ }
+
// ALTER_PARTITION_STATES replay applies a detached snapshot here,
mirroring alterIvmInfo(). Live
// invalidation changes submit their journal from the mutating method
instead.
//
// A payload without the member carries no state at all, which is not the
same as an empty map that
// says the states are now empty: leaving them alone is the only answer
that cannot lose state.
public void alterPartitionStates(Map<String, MTMVPartitionState>
partitionStates) {
- if (partitionStates == null) {
+ replayAlterPartitionStates(partitionStates, null, false);
+ }
+
+ /**
+ * ALTER_PARTITION_STATES replay: applies the states the payload carries,
and drops the snapshots it
+ * names. Both in one lock acquisition, because a reader that saw the new
requirement while the
+ * snapshot was still there could let a transparent rewrite serve rows the
rebuild has to replace.
+ *
+ * <p>A payload without the states carries none, which is not the same as
an empty map that says the
+ * states are now empty: leaving them alone is the only answer that cannot
lose state.
+ *
+ * <p>{@code merge} says what the payload's states are. False, which is
what a payload written before the
+ * member existed means and what the invalidation channel still writes,
carries the map itself and
+ * replaces. True carries only the partitions a change touched -- see
submitPartitionStatesDelta -- so the
+ * entries it does not name belong to other records (an invalidation that
ran during the refresh, an entry
+ * alignment added) and are merged over rather than dropped.
+ */
+ public void replayAlterPartitionStates(Map<String, MTMVPartitionState>
partitionStates,
+ Set<String> removedSnapshotPartitions, boolean merge) {
+ writeMvLock();
+ try {
+ if (partitionStates != null) {
+ if (merge) {
+ for (Entry<String, MTMVPartitionState> entry :
partitionStates.entrySet()) {
+ this.partitionStates.put(entry.getKey(), new
MTMVPartitionState(entry.getValue()));
+ }
+ } else {
+ this.partitionStates =
MTMVPartitionState.copyOf(partitionStates);
+ }
+ }
+ refreshSnapshot.removeSnapshots(removedSnapshotPartitions);
+ } finally {
+ writeMvUnlock();
+ }
+ }
+
+ /**
+ * The {@code latestEpoch} of the given MV partitions, taken under the MV
read lock.
+ *
+ * <p>This is the value a refresh has to remember: what it read from the
base tables is described by
+ * the requirement in force when it started reading, so writing that value
back as the new
+ * {@code refreshEpoch} is what keeps an invalidation arriving mid-refresh
from being swallowed. A
+ * partition without an entry is left out -- a caller writes an epoch only
for what it captured.
+ */
+ public Map<String, Long> getLatestEpochs(Set<String> partitionNames) {
+ if (CollectionUtils.isEmpty(partitionNames)) {
+ return Collections.emptyMap();
+ }
+ // Sized before the lock: the state map is what needs it, and building
the map is not part of that.
+ Map<String, Long> res =
Maps.newHashMapWithExpectedSize(partitionNames.size());
+ readMvLock();
+ try {
+ for (String partitionName : partitionNames) {
+ MTMVPartitionState state = partitionStates.get(partitionName);
+ if (state != null) {
+ res.put(partitionName, state.getLatestEpoch());
+ }
+ }
+ return res;
+ } finally {
+ readMvUnlock();
+ }
+ }
+
+ /**
+ * Brings the partition states in line with the MV's partitions: every
partition gets an entry, and
+ * every entry whose partition is gone is dropped.
+ *
+ * <p>Alignment is what makes "the partition exists" and "the entry
exists" the same thing, and it is
+ * why an invalidation cannot miss: rows are only written by a refresh,
and every refresh aligns
+ * before it reads a base table, so a partition that holds rows always has
an entry for the mark to
+ * land on. The other direction is what makes the criterion safe -- an
entry created here describes a
+ * partition with no rows yet, so requiring one generation of it discards
no requirement that was
+ * made earlier.
+ *
+ * <p>What it changes is journaled, because the entry has to be on disk
before the rows it describes
+ * can be: a crash between this and the task result would otherwise leave
a partition that holds rows
+ * with no entry at all, and every later invalidation of it would find
nothing to land on. That is the
+ * one shape in which the criterion cannot be read -- "no entry" is
supposed to mean "no rows" -- so
+ * the entry is made durable before any base table is read rather than
derived again on the next run.
+ *
+ * <p>It is deliberately not a hook on every path that creates or drops a
partition. An entry is
+ * derived state, and rebuilding it from the live partition set also
repairs whatever a crash left
+ * behind: the drop of a partition and the removal of its entry are two
journal records, and only
+ * their order -- partition first -- is safe, which leaves at most a stale
entry that the next
+ * alignment drops.
+ *
+ * <p>Only an IVM MV is aligned. For a non-IVM MV the map stays as it is,
and every reader treats
+ * "empty" and "no state" the same.
+ */
+ public void alignPartitionStates() {
+ if (!isIvm()) {
return;
}
+ EditLogItem editLogItem = null;
writeMvLock();
try {
- this.partitionStates = MTMVPartitionState.copyOf(partitionStates);
+ // Read here rather than handed in by the caller: a caller has to
read the names before it takes
+ // this lock, and a partition created in between -- by a
concurrent refresh's partition sync --
+ // would then be dropped by the retainAll below, taking with it
the state a following
+ // invalidation has to land on. The read is cheap and takes no
lock of its own, so doing it here
+ // does not add an edge to the lock order.
+ Set<String> livePartitions = Sets.newHashSet(getPartitionNames());
+ boolean changed =
partitionStates.keySet().retainAll(livePartitions);
+ for (String partitionName : livePartitions) {
+ if (!partitionStates.containsKey(partitionName)) {
+ partitionStates.put(partitionName,
MTMVPartitionState.initial());
+ changed = true;
+ }
+ }
+ if (changed) {
+ editLogItem =
submitPartitionStatesChange(Collections.emptySet());
+ }
} finally {
writeMvUnlock();
}
+ if (editLogItem != null) {
+ editLogItem.await();
+ }
}
- public void invalidateIvmBaseline() {
- EditLogItem editLogItem;
+ /**
+ * The snapshots of the partitions that are clean after this result's
epochs were applied.
+ *
+ * <p>An invalidation that reached a partition while the task ran leaves
it dirty, and its snapshot
+ * must stay gone: dropping the entry is what keeps transparent rewrite
away from rows the rebuild has
+ * to replace, and a result written back afterwards would undo exactly
that. Removing only the entry
+ * keeps the rest of the map, which the removal on the invalidation side
cannot express.
+ *
+ * <p>The caller holds the MV write lock and has already applied the
epochs, so {@code isDirty} here
+ * reads the state the data is actually described by.
+ *
+ * <p>Only an IVM MV has partition states, so only its write-back is
narrowed here: every entry of a
+ * non-IVM MV has no state to be dirty in and is written back as it always
was.
+ */
+ private Map<String, MTMVRefreshPartitionSnapshot>
snapshotsOfCleanPartitions(
+ Map<String, MTMVRefreshPartitionSnapshot> snapshots) {
+ if (MapUtils.isEmpty(snapshots)) {
+ return snapshots;
Review Comment:
Confirmed, and it is the empty case that leaks the alias -- a non-empty map
goes through the new map the filter builds. `after()` hands over the task's own
`partitionSnapshots` field, which is the map the worker merges every committed
batch into, and `executeCancelLogic` calls `after()` from the cancel thread
(`cancel(false)`, so the execution goes on). The record is written out
asynchronously, so what serializes is whatever the map holds by then.
`42ed8c44575` closes it in both places: the empty payload is an immutable
empty map, and `after()` hands over a copy rather than the map itself -- the
rule `getIvmCapturedEpochs` already follows for the epochs, for the same
reason. The copy also covers what the filter does not run on: a plain MV's
result, and a non-empty payload the worker keeps adding to.
`MTMVTest#testAnEmptySnapshotPayloadIsNotTheMapTheTaskKeepsFilling` builds a
result from an empty map, adds an entry afterwards as the worker would, and
asserts the journaled payload is still empty;
`MTMVTaskTest#testTheSnapshotsHandedOverAreNotTheMapTheWorkerKeepsFilling` pins
the handover itself.
--
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]