yujun777 commented on code in PR #68390:
URL: https://github.com/apache/doris/pull/68390#discussion_r4119929202
##########
fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java:
##########
@@ -947,67 +1297,118 @@ private MTMVRelatedTableIf findPctTable(BaseTableInfo
baseTableInfo) {
return null;
}
+ private EditLogItem submitIvmInfoChange() {
+ // The caller has already mutated ivmInfo under the MV write lock.
Submit its snapshot directly;
+ // replay later applies the payload through alterIvmInfo().
+ AlterMTMV alterMTMV = new AlterMTMV(
+ new TableNameInfo(getQualifiedDbName(), getName()),
MTMVAlterOpType.ALTER_IVM_INFO);
+ alterMTMV.setIvmInfo(ivmInfo);
+ return submitAlterLog(alterMTMV);
+ }
+
/**
- * Release the IVM baseline barrier after the partitions it named have
been rebuilt, or after
- * partition sync removed them (a dropped partition resolves its own
entry: the partition and its
- * IVM offsets are both gone).
+ * Raises the requirement of the given MV partitions, so the next refresh
rebuilds them.
*
- * <p>Guarded by schemaChangeVersion, like {@link
#persistIvmBaselineGuard}: a base-table change
- * landing while the rebuild runs carries its own barrier entry, and a
blind clear would swallow
- * it. Failing instead preserves that entry -- the next refresh rebuilds
it together with the
- * partitions this task handled.
- *
- * <p>Journals the new state right away, like every other ivmInfo mutation
here. A task that dies
- * before {@link #addTaskResult} would otherwise leave the release in
memory only, and a restart
- * would resurrect the barrier from disk.
+ * <p>This is an invalidation-shaped mutation, journaled as the whole
state map before whatever needs
+ * it is done. A caller about to make a partition's rows unusable says so
with it: the raised
+ * requirement survives a crash, so a refresh that never got to publish
its rebuild leaves partitions
+ * naming a generation they do not hold, and the next refresh rebuilds
them.
*/
- public void releaseIvmBaselineRebuild(long expectedSchemaChangeVersion)
throws JobException {
+ public void markPartitionsForRebuild(Set<String> partitionNames) {
+ if (CollectionUtils.isEmpty(partitionNames)) {
+ return;
+ }
EditLogItem editLogItem;
writeMvLock();
try {
- if (ivmInfo == null || !ivmInfo.isBaselineRebuildRequired()) {
- // Nothing to release: skip both the mutation and the journal
entry. Any base-table
- // change that raced us in is still caught by
validateIvmRefreshStart() below.
- return;
+ boolean changed = false;
+ for (String partitionName : partitionNames) {
+ MTMVPartitionState state = partitionStates.get(partitionName);
+ if (state == null) {
+ // A partition dropped since the caller planned it has no
rows to protect.
+ continue;
+ }
+ state.setLatestEpoch(state.getLatestEpoch() + 1);
+ changed = true;
}
- if (schemaChangeVersion != expectedSchemaChangeVersion) {
- throw new JobException("Base table metadata changed before IVM
baseline refresh, mv="
- + getName());
+ if (!changed) {
+ return;
}
- ivmInfo.clearBaselineRebuild();
- editLogItem = submitIvmInfoChange();
+ editLogItem = submitPartitionStatesChange(Collections.emptySet());
} finally {
writeMvUnlock();
}
editLogItem.await();
}
- public void persistIvmBaselineGuard(RefreshMode refreshMode, Set<String>
baselinePartitions,
- long expectedSchemaChangeVersion) throws JobException {
+ /**
+ * Raises the requirement of the given MV partitions that do not name one,
and reports what each of those
+ * partitions now names.
+ *
+ * <p>This is what a refresh about to replace a partition says about it,
and it has to be said before that
+ * replacement reads anything: an overwrite is two halves -- the rows are
committed into temporary
+ * partitions, and a swap publishes them -- so a refresh that dies in
between leaves the live partition
+ * holding the rows it had while whatever its read consumed, the offsets
of the streams it read among
+ * them, is already committed with the first half. The epochs a refresh
records ride with its result, and
+ * a refresh that never returns records none, so without this nothing
would say the partition owes the
+ * rebuild and the next refresh would read on from an offset past a change
the partition never received.
+ *
+ * <p>What the caller gets back is the requirement it raised, which is the
ceiling its write-back is
+ * clamped to; see MTMVTask's captured epochs. A partition that already
names a requirement is left alone
+ * and is not part of that result: it names the requirement the refresh
answers for, and the caller must
+ * record what it read rather than what it found. Raising it again would
move it above that, and the
+ * partition would be rebuilt a second time for nothing. The record it
submits still carries the whole
+ * map -- that is what this channel carries -- but it is submitted only
when something was raised.
+ *
+ * <p>This differs from {@link #markPartitionsForRebuild} on purpose: that
one is an invalidation, and it
+ * raises the requirement of every partition it names because it has to
outrank a refresh already
+ * running. This one is a refresh's own record of what it is about to do,
and a partition that already
+ * names such a requirement does not need a second one.
+ *
+ * <p>Which partitions need it is decided under the same lock as the
raise. Read outside it, a mark
+ * landing in between would leave this call blind to a requirement it then
raises above, and the caller
+ * would record its own value as met for a change that arrived after it
read.
+ */
+ public Map<String, Long> raiseRebuildRequirement(Set<String>
partitionNames) {
+ if (CollectionUtils.isEmpty(partitionNames)) {
+ return Collections.emptyMap();
+ }
+ Map<String, Long> raised =
Maps.newHashMapWithExpectedSize(partitionNames.size());
EditLogItem editLogItem;
writeMvLock();
try {
- if (schemaChangeVersion != expectedSchemaChangeVersion) {
- throw new JobException("Base table metadata changed before IVM
baseline refresh, mv=" + getName());
+ for (String partitionName : partitionNames) {
+ MTMVPartitionState state = partitionStates.get(partitionName);
+ if (state == null || state.isDirty()) {
+ // Dropped since the caller planned it, or already naming
a requirement of its own.
+ continue;
+ }
+ state.setLatestEpoch(state.getLatestEpoch() + 1);
+ raised.put(partitionName, state.getLatestEpoch());
}
- if (refreshMode == RefreshMode.COMPLETE) {
- ivmInfo.requireCompleteBaselineRebuild();
- } else {
-
ivmInfo.addPendingBaselineRebuildPartitions(baselinePartitions);
+ if (raised.isEmpty()) {
+ return Collections.emptyMap();
}
- editLogItem = submitIvmInfoChange();
+ editLogItem = submitPartitionStatesChange(Collections.emptySet());
Review Comment:
Confirmed, and measured: on a thousand-partition MV, a one-partition raise
was writing a record with 1000 entries under the MV write lock. `616714c79f0`
makes that record the delta it is, and the test in it fails with `expected: <1>
but was: <1000>` when the payload goes back to the whole map.
**What the record carries now.** `raiseRebuildRequirement` submits the
entries it raised and only those. `AlterMTMV` gets `mergePartitionStates`
(`mps`), and `MTMV#replayAlterPartitionStates` merges when it is set instead of
replacing the map. A payload written before the member existed carries the
whole map and an absent member reads as `false`, so every record already in a
journal -- and the invalidation channel -- keep replacing exactly as before.
The delta's entries are detached copies, so the payload still describes the
raise rather than the MV's live entry, which is the property the earlier round
asked for.
**Why merging is safe here.** The replay applies records in the order the MV
write lock fixed when they were enqueued. A later whole-map record was copied
from the map a merge had already produced, so it carries the merged entries
too; and an alignment that drops entries takes them out in a record that is
replayed after the deltas that named them. The entries a delta does not name
are the ones that belong to other records -- an invalidation that ran during
the refresh, an entry an alignment added -- which is exactly what replacing
with the delta would have dropped. The ADD_TASK channel already journals a
delta and merges it on replay for the same reason; this makes the refresh's
pre-overwrite record use that shape.
**What I did not change, deliberately.** `markPartitionsForRebuild`,
`invalidateIvmBaseline` and `alignPartitionStates` still journal the whole map.
Your finding is about the pre-overwrite record, which runs on every
partition-based refresh; those three run on a DDL or an alignment, and their
record is the map itself -- which is also what lets one of them carry a map
that lost entries. If you want them moved to the delta shape as well, that is a
separate change, and it needs `alignPartitionStates` to keep a way to remove
entries.
Covered by
`MTMVTest#testRaiseRebuildRequirementJournalsTheDeltaAndReplaysItMerged`: a
thousand-partition MV, one partition raised, the payload carries one entry, it
survives the journal as itself, a replay merges it over the states the MV
holds, and a whole-map record still replaces. 143 unit tests across `MTMVTest`,
`MTMVTaskTest` and `AlterMTMVTest`, checkstyle clean, and the regression suites
that read the epochs (`test_ivm_overwrite_failure_between_the_halves`,
`test_ivm_partition_epoch_rebuild`,
`test_ivm_partitions_after_failed_complete`, `test_ivm_baseline_marker_scope`)
are green.
--
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]