yujun777 commented on code in PR #68390:
URL: https://github.com/apache/doris/pull/68390#discussion_r4083026833
##########
fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java:
##########
@@ -715,13 +685,71 @@ 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().needsRebuild()) {
+ dirtyPartitions.add(plannedState.getKey());
+ }
+ }
+ this.ivmPlannedEpochs = plannedEpochs;
+ Map<String, MTMVRefreshPartitionSnapshot> rebuiltSnapshots =
Maps.newHashMap();
+ List<String> rebuildScope = Lists.newArrayList();
+ Set<String> rebuildCompleted = Sets.newLinkedHashSet();
+ 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.needRefreshPartitions = toRebuild;
+ this.refreshMode = generateRefreshMode(toRebuild);
+ try {
+ executePartitionBasedRefresh(refreshContext,
RefreshMode.PARTITIONS, ctx);
+ } finally {
+ // Counted from the groups that committed, not from the ones
that were planned: a refresh
+ // that failed part-way through the rebuild must not report
partitions it never replaced.
+ recordRebuiltPartitions(request, partitionSnapshots.size());
+ }
+ rebuiltSnapshots.putAll(partitionSnapshots);
+ // Kept before the incremental attempt resets the accumulators to
its own scope: both phases
+ // belong to this refresh, so the progress it reports is the union
of the two.
+ rebuildScope.addAll(needRefreshPartitions);
+ rebuildCompleted.addAll(completedPartitions);
+ }
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 incremental attempt reset the accumulators it owns, so
the rebuild's are merged back
+ // here: its batches committed, and without them the
partitions it rebuilt would look
+ // unsynced and be refreshed again on every following round.
+ this.partitionSnapshots.putAll(rebuiltSnapshots);
+ // The incremental attempt reset the accumulators to its own
scope. Both phases are part of
+ // the refresh that is being reported, so the denominator is
the union of the two and the
+ // completed side keeps what each phase committed: a refresh
that rebuilt one partition and
+ // caught up another would otherwise record two of one.
+ Set<String> mergedScope = Sets.newLinkedHashSet(rebuildScope);
+ mergedScope.addAll(needRefreshPartitions);
+ this.needRefreshPartitions = Lists.newArrayList(mergedScope);
+ this.completedPartitions.addAll(rebuildCompleted);
return AttemptResultType.SUCCESS;
}
if (ivmResult.getFailureReason() !=
IvmFailureReason.MV_PARTITION_NOT_FOUND) {
Review Comment:
Implemented, and then simplified. The task keeps one snapshot accumulator
instead of resetting it per phase, so what a rebuild committed is still in it
when the fallback plans, and the plan subtracts the partitions already there
(`MTMVTask#excludingRebuiltPartitions`). The second map this originally needed
is gone in `ee680f0e6a3`.
##########
fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelationManager.java:
##########
@@ -359,16 +362,48 @@ public void alterTable(BaseTableInfo oldTableInfo,
Optional<BaseTableInfo> newTa
// when replace, need deal two table
if (isReplace) {
// REPLACE TABLE already invalidates the IVM baseline explicitly,
see Alter#processReplaceTable
- processBaseTableChange(newTableInfo.get(), "The base table has
been updated:", false);
+ processBaseTableChange(newTableInfo.get(), "The base table has
been updated:", false, false);
}
- // A RENAME leaves every column alone, and the failure it does cause
-- the MV query still
- // spells the old name -- is already reported by the refresh itself
(MTMVTask#run resolves
- // the base tables from the query before it ever looks at the
baseline). Invalidating here
- // would only leave a stale flag behind: rename the table back and the
query is analyzable
- // again, yet every strict INCREMENTAL refresh would stay rejected
until a COMPLETE one ran.
boolean renamed = !isReplace && newTableInfo.isPresent()
&& !Objects.equals(oldTableInfo.getTableName(),
newTableInfo.get().getTableName());
- processBaseTableChange(oldTableInfo, "The base table has been
updated:", !renamed);
+ // The invalidation runs first, while the dependencies are still
registered under the name the
+ // rename is leaving: moving them first would make this lookup --
which is by the old name -- find
+ // nothing, and the rename would stop invalidating anything at all.
+ processBaseTableChange(oldTableInfo, "The base table has been
updated:", !renamed, renamed);
+ if (renamed) {
+ renameBaseTable(oldTableInfo, newTableInfo.get());
+ }
+ }
+
+ /**
+ * Move a renamed table's entries in the dependency maps to its new name.
+ *
+ * <p>The maps are keyed by {@link BaseTableInfo}, which compares by name,
and an MV keeps the relation
+ * it was created against -- a rename leaves the MV query spelling the old
name, so it no longer
+ * analyzes and the relation is not recomputed. Without this the maps
would keep the old name, and a
+ * metadata-only change to the table under its new name -- a TRUNCATE,
say, which emits no row binlog --
+ * would find no dependent MV to invalidate. Renaming the table back then
restores an analyzable query
+ * whose MV still holds the rows that change removed, and nothing names
the partition that would have
+ * to be rebuilt. Moving the entries is what a rename needs instead of the
invalidation it used to
+ * carry: a rename changes no rows, so there is nothing to rebuild, only a
lookup that has to keep
+ * working.
+ */
+ private void renameBaseTable(BaseTableInfo oldTableInfo, BaseTableInfo
newTableInfo) {
+ moveRelationKey(tableMTMVs, oldTableInfo, newTableInfo);
+ moveRelationKey(tableMTMVsOneLevelAndFromView, oldTableInfo,
newTableInfo);
+ }
+
+ private void moveRelationKey(Map<BaseTableInfo, Set<BaseTableInfo>> map,
+ BaseTableInfo oldTableInfo, BaseTableInfo newTableInfo) {
+ Set<BaseTableInfo> dependents = map.get(oldTableInfo);
+ if (CollectionUtils.isEmpty(dependents)) {
+ return;
+ }
+ // Registered under the new name before the old one is dropped: a
concurrent base-table change
Review Comment:
Dropped with the re-key: the transition is not made, so there is nothing to
race. A rename invalidates every MV that reads the table, which is what it did
before this PR, and the cost is one extra COMPLETE for a rename-and-rename-back.
--
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]