github-actions[bot] commented on code in PR #68390:
URL: https://github.com/apache/doris/pull/68390#discussion_r4122444027
##########
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:
[P1] Redeem held rebuild records only after a delta actually runs. A partial
rebuild of dirty `p` can read non-PCT `D` at its older stream offset and hold
`p`'s epoch/current snapshot. If a newly populated base partition makes the
first delta return `MV_PARTITION_NOT_FOUND`, retry sync can replace the only
clean MV partition with a newly dirty one. The recalculated `incrementalScope`
then contains only dirty names, so `executeSingleIvmAttempt` returns success
without calling `doRefresh`; this line still redeems `p`. `p` can retain rows
from `D=10` while its published snapshot says `D=20`, making stale rows appear
current to rewrite. The prior :1283 thread covers a successful delta that
*ran*; this path runs none. Keep held records dirty until a physical delta
succeeds, and test the retry-to-empty sequence.
##########
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:
[P2] Keep strict `INCREMENTAL` out of the schema-change COMPLETE shortcut.
Manual `REFRESH ... INCREMENTAL` defaults to `allowFallback=false`, but this
condition rejects only strict `PARTITIONS`; the branch below returns
`[COMPLETE]`, reconciles an unusable stream, and performs a full RESET/SNAPSHOT
rebuild. Thus a strict request can succeed after the very stream reset that the
all-dirty shortcut guards with `allowFallback`. The existing :527 thread
addressed strict `PARTITIONS` scope, and :584 addressed a different all-dirty
branch; this is the still-unguarded strict `INCREMENTAL` route. Refuse it with
an actionable COMPLETE/AUTO or explicit FALLBACK message, and cover
SCHEMA_CHANGE with an unusable stream.
##########
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:
[P2] Journal only the marked state entries for a partition-scoped base DDL
invalidation. For a TRUNCATE affecting one PCT partition of an MV with `N`
partitions, `markIvmPartitionsInvalidated` changes one state and removes its
snapshot, but this call passes the entire state map to
`AlterMTMV.setPartitionStates`, copying and journaling `N` entries under
`mvRwLock`. The :1392 thread fixed the separate pre-overwrite raise using a
merge delta; this DDL caller still scales with unrelated partitions and never
deletes state entries. Persist a merge delta for `marked` together with the
named snapshot removals, and add a one-partition journal-size test on a large
MV.
##########
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:
[P2] Use a delta when alignment only adds a state entry. A refresh sync that
adds one MV partition leaves `retainAll` unchanged, inserts one `{0,1}` state,
then reaches this call, which deep-copies and journals all `N+1` partition
states under the MV lock. This differs from the DDL invalidation caller and the
fixed pre-overwrite raise: alignment may also delete entries, so full
replacement is needed only in that case. Track additions and removals
separately; merge just the new entries on an add-only alignment, and retain
full replacement when an old entry was removed. A large-MV one-addition payload
test would catch this.
##########
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:
[P1] Detach the empty snapshot map before building the queued `ADD_TASK`
record. STOP uses `cancel(false)`, so `after()` can pass the worker's
still-mutable `partitionSnapshots` here while it is empty. The leader's
`updateSnapshots` applies no entries, but `AlterMTMV.setPartitionSnapshots`
retains this same map and batch edit-log submission queues the object. If the
worker adds a snapshot before serialization, replay installs an entry the
leader never applied, possibly after a partition invalidation removed it. The
existing :271 thread covers cancellation and detached epoch captures; this new
empty-map alias is a separate snapshot/replay mismatch. Return an
immutable/detached empty map and test leader/replay equality across the
cancellation and queue window.
```suggestion
return Collections.emptyMap();
```
--
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]