github-actions[bot] commented on code in PR #68390:
URL: https://github.com/apache/doris/pull/68390#discussion_r4119679593
##########
fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java:
##########
@@ -550,14 +676,43 @@ private PartitionRefreshPlan
planPartitionRefresh(MTMVRefreshContext context,
+ "does not support refreshing by partition");
}
try {
+ Set<String> planned =
Sets.newLinkedHashSet(MTMVPartitionUtil.getMTMVNeedRefreshPartitions(context,
+ relation.getBaseTablesOneLevelAndFromView()));
+ planned.addAll(rebuildRequired);
return PartitionRefreshPlan.success(context,
- MTMVPartitionUtil.getMTMVNeedRefreshPartitions(context,
- relation.getBaseTablesOneLevelAndFromView()));
+ excludingRebuiltPartitions(Lists.newArrayList(planned)));
} catch (Exception e) {
return PartitionRefreshPlan.fallback(e.getMessage());
}
}
+ /**
+ * Takes the partitions this task has already replaced out of a planned
set.
+ *
+ * <p>The plan is computed from the snapshot the MV holds, which this task
has not published yet, so a
+ * partition this task has just filled still looks unsynced to it.
Refreshing it here would replace the
+ * rows it holds with a second read of the same base table, and would do
it in the one case that has
+ * already paid for it: the same refresh falling back out of the
incremental attempt. What the
+ * accumulator holds at this point is exactly those partitions -- the
fallback runs after an attempt
+ * that committed nothing, and a plan is built before anything is written.
+ *
+ * <p>An explicit partition list is left alone. It is the request itself
rather than an inference from
+ * the MV's snapshot, and the rebuild phase does not take partitions out
of it either: what the request
+ * names is refreshed, and at worst it is refreshed twice within one task.
+ */
+ private List<String> excludingRebuiltPartitions(List<String>
plannedPartitions) {
+ if (partitionSnapshots.isEmpty()) {
+ return plannedPartitions;
+ }
+ List<String> remaining =
Lists.newArrayListWithCapacity(plannedPartitions.size());
+ for (String partitionName : plannedPartitions) {
+ if (!partitionSnapshots.containsKey(partitionName)) {
Review Comment:
[P1] Preserve a current non-PCT baseline when excluding rebuilt partitions
from fallback. Suppose a partitioned F JOIN D aggregate has p1/p2 and D changes
10 to 20 while p1 is dirty. The pre-rebuild of p1 reads D SNAPSHOT at the old
offset (10). If the IVM delta returns a non-COMPLETE fallback reason, this
helper removes p1 from the PARTITIONS plan, leaving only p2; that partial
overwrite also reads D SNAPSHOT at 10, so the task can report SUCCESS with both
partitions at 10 while recording D's current version (20) for them. Later sync
checks can then skip the still-pending D delta and transparent rewrite can
serve wrong rows. The old full-scope fallback reset D and repaired both
partitions; the earlier :815 thread covered its repeated overwrite, not this
new stale-success outcome. Reconcile D's pending stream change or take a full
current-snapshot rebuild before publishing success, and test exact rows after a
forced non-COMPLETE IVM fallback.
##########
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:
[P2] Persist only the selected partition's new rebuild requirement here.
Each IVM PARTITIONS refresh of one clean partition calls this method before
overwrite, then submitPartitionStatesChange copies every partition state under
mvRwLock and serializes the whole map to the edit log. On an MV with many
partitions, refreshing just one becomes a large lock-held copy and journal
record each time this path runs. The earlier ADD_TASK thread addressed that
later task-result payload; this new routine pre-overwrite record has the same
scaling problem. Use a merge-on-replay epoch delta for the selected entries and
cover one-partition journal size on a large MV.
##########
fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java:
##########
@@ -715,12 +768,47 @@ 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);
Review Comment:
[P1] Do not publish this rebuild as fresh if the following delta fails. A
strict INCREMENTAL can rebuild dirty p1 using a non-PCT table's SNAPSHOT at its
old stream offset, while generatePartitionSnapshots records that table's
current version. For a partitioned F JOIN D aggregate with D changing 10 to 20,
p1's overwrite still holds 10; if the subsequent IVM delta returns a fallback
reason, strict mode fails, but onFail publishes p1's captured epoch and current
D snapshot. FAILED MVs remain eligible for transparent rewrite, and the
snapshot comparisons then admit p1 although the base query returns 20. The old
strict path rejected this pending rebuild. Keep such partitions
dirty/unavailable until the required delta commits, or tag the rebuild with the
historical source snapshot; add a forced-delta-failure regression that checks
rewrite eligibility and rows.
--
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]