MTMVTask.java
// Licensed to the Apache Software Foundation (ASF) under one
// or more contributor license agreements. See the NOTICE file
// distributed with this work for additional information
// regarding copyright ownership. The ASF licenses this file
// to you under the Apache License, Version 2.0 (the
// "License"); you may not use this file except in compliance
// with the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.
package org.apache.doris.job.extensions.mtmv;
import org.apache.doris.analysis.PartitionKeyDesc;
import org.apache.doris.catalog.Column;
import org.apache.doris.catalog.Database;
import org.apache.doris.catalog.DatabaseIf;
import org.apache.doris.catalog.Env;
import org.apache.doris.catalog.MTMV;
import org.apache.doris.catalog.OlapTable;
import org.apache.doris.catalog.Partition;
import org.apache.doris.catalog.ScalarType;
import org.apache.doris.catalog.TableIf;
import org.apache.doris.catalog.info.TableNameInfo;
import org.apache.doris.catalog.stream.BaseTableStream;
import org.apache.doris.catalog.stream.StreamReadMode;
import org.apache.doris.cloud.qe.ComputeGroupException;
import org.apache.doris.common.AnalysisException;
import org.apache.doris.common.Config;
import org.apache.doris.common.DdlException;
import org.apache.doris.common.FeConstants;
import org.apache.doris.common.Pair;
import org.apache.doris.common.Status;
import org.apache.doris.common.UserException;
import org.apache.doris.common.util.DebugPointUtil;
import org.apache.doris.common.util.DebugUtil;
import org.apache.doris.common.util.MetaLockUtils;
import org.apache.doris.common.util.TimeUtils;
import org.apache.doris.common.util.Util;
import org.apache.doris.datasource.CatalogIf;
import org.apache.doris.datasource.mvcc.MvccSnapshot;
import org.apache.doris.datasource.mvcc.MvccTable;
import org.apache.doris.datasource.mvcc.MvccTableInfo;
import org.apache.doris.job.common.TaskStatus;
import org.apache.doris.job.exception.JobException;
import org.apache.doris.job.task.AbstractTask;
import org.apache.doris.metric.MetricRepo;
import org.apache.doris.mtmv.BaseColInfo;
import org.apache.doris.mtmv.BaseTableInfo;
import org.apache.doris.mtmv.MTMVBaseTableIf;
import org.apache.doris.mtmv.MTMVPartitionInfo.MTMVPartitionType;
import org.apache.doris.mtmv.MTMVPartitionState;
import org.apache.doris.mtmv.MTMVPartitionUtil;
import org.apache.doris.mtmv.MTMVPlanUtil;
import org.apache.doris.mtmv.MTMVRefreshContext;
import org.apache.doris.mtmv.MTMVRefreshEnum.MTMVRefreshState;
import org.apache.doris.mtmv.MTMVRefreshEnum.MTMVState;
import org.apache.doris.mtmv.MTMVRefreshEnum.RefreshMethod;
import org.apache.doris.mtmv.MTMVRefreshPartitionSnapshot;
import org.apache.doris.mtmv.MTMVRelatedTableIf;
import org.apache.doris.mtmv.MTMVRelation;
import org.apache.doris.mtmv.MTMVUtil;
import org.apache.doris.mtmv.ivm.IvmFailureReason;
import org.apache.doris.mtmv.ivm.IvmIncrRefreshContext;
import org.apache.doris.mtmv.ivm.IvmIncrRefreshManager;
import org.apache.doris.mtmv.ivm.IvmIncrRefreshResult;
import org.apache.doris.mtmv.ivm.IvmPlanSignature;
import org.apache.doris.mtmv.ivm.IvmRewriteContext;
import org.apache.doris.mtmv.ivm.IvmUtil;
import org.apache.doris.nereids.NereidsPlanner;
import org.apache.doris.nereids.StatementContext;
import org.apache.doris.nereids.trees.plans.commands.CreateMTMVCommand;
import org.apache.doris.nereids.trees.plans.commands.UpdateMvByPartitionCommand;
import org.apache.doris.nereids.trees.plans.commands.info.RefreshMTMVInfo.RefreshMode;
import org.apache.doris.qe.ConnectContext;
import org.apache.doris.qe.OriginStatement;
import org.apache.doris.qe.QeProcessorImpl;
import org.apache.doris.qe.StmtExecutor;
import org.apache.doris.rpc.RpcException;
import org.apache.doris.system.SystemInfoService;
import org.apache.doris.thrift.TCell;
import org.apache.doris.thrift.TRow;
import org.apache.doris.thrift.TStatusCode;
import com.google.common.base.Strings;
import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableMap;
import com.google.common.collect.Lists;
import com.google.common.collect.Maps;
import com.google.common.collect.Sets;
import com.google.gson.Gson;
import com.google.gson.annotations.SerializedName;
import org.apache.commons.collections4.CollectionUtils;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.util.Collection;
import java.util.Collections;
import java.util.Comparator;
import java.util.List;
import java.util.Map;
import java.util.Map.Entry;
import java.util.Objects;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.Callable;
import java.util.concurrent.ConcurrentSkipListSet;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Consumer;
public class MTMVTask extends AbstractTask {
private static final Logger LOG = LogManager.getLogger(MTMVTask.class);
public static final int DEFAULT_REFRESH_PARTITION_NUM = 1;
public static final String DEBUG_POINT_SKIP_PARTITION_SYNC =
"MTMVTask.syncPartitionsIfNeeded.skip";
public static final String DEBUG_POINT_SKIP_PARTITION_SYNC_FILTER =
"MTMVTask.syncPartitionsIfNeeded.skip.filter";
private static final Gson GSON = new Gson();
public static final ImmutableList<Column> SCHEMA = ImmutableList.of(
new Column("TaskId", ScalarType.createStringType()),
new Column("JobId", ScalarType.createStringType()),
new Column("JobName", ScalarType.createStringType()),
new Column("MvId", ScalarType.createStringType()),
new Column("MvName", ScalarType.createStringType()),
new Column("MvDatabaseId", ScalarType.createStringType()),
new Column("MvDatabaseName", ScalarType.createStringType()),
new Column("Status", ScalarType.createStringType()),
new Column("ErrorMsg", ScalarType.createStringType()),
new Column("CreateTime", ScalarType.createStringType()),
new Column("StartTime", ScalarType.createStringType()),
new Column("FinishTime", ScalarType.createStringType()),
new Column("DurationMs", ScalarType.createStringType()),
new Column("TaskContext", ScalarType.createStringType()),
new Column("RefreshMode", ScalarType.createStringType()),
new Column("NeedRefreshPartitions", ScalarType.createStringType()),
new Column("CompletedPartitions", ScalarType.createStringType()),
new Column("Progress", ScalarType.createStringType()),
new Column("LastQueryId", ScalarType.createStringType()),
new Column("ComputeGroup", ScalarType.createStringType()),
new Column("IvmFallbackReason", ScalarType.createStringType()),
new Column("IvmRebuiltPartitions", ScalarType.createStringType()));
public static final ImmutableMap<String, Integer> COLUMN_TO_INDEX;
static {
ImmutableMap.Builder<String, Integer> builder = new ImmutableMap.Builder<String, Integer>();
for (int i = 0; i < SCHEMA.size(); i++) {
builder.put(SCHEMA.get(i).getName().toLowerCase(), i);
}
COLUMN_TO_INDEX = builder.build();
}
public enum MTMVTaskTriggerMode {
MANUAL,
COMMIT,
SYSTEM
}
public enum MTMVTaskRefreshMode {
COMPLETE,
PARTIAL,
NOT_REFRESH
}
private enum RefreshAttemptType {
IVM,
PARTITIONS,
COMPLETE
}
private enum AttemptResultType {
SUCCESS,
// The current attempt failed before writing MV data, so the task may
// continue to the next configured fallback attempt.
FALLBACK_ALLOWED,
// A previous IVM delta may have partially written data. PARTITIONS
// cannot prove it repairs that state, so recovery must be COMPLETE.
FALLBACK_TO_COMPLETE
}
private static class RefreshRequest {
private final RefreshMode refreshMode;
private final boolean allowFallback;
private final List<String> partitions;
// True only for REFRESH ... PARTITION(S). Explicit partition refresh is
// a user-selected scope and must not expand to COMPLETE via fallback.
private final boolean explicitPartitions;
private RefreshRequest(RefreshMode refreshMode, boolean allowFallback,
List<String> partitions, boolean explicitPartitions) {
this.refreshMode = Objects.requireNonNull(refreshMode, "refreshMode can not be null");
this.allowFallback = allowFallback;
this.partitions = partitions == null ? Lists.newArrayList() : partitions;
this.explicitPartitions = explicitPartitions;
}
}
private static class PartitionRefreshPlan {
private final MTMVRefreshContext context;
// False means partition planning failed before any refresh write. The
// caller may convert it to COMPLETE only when the request allows fallback.
private final boolean canRefreshByPartitions;
private final List<String> partitions;
private final String fallbackReason;
private PartitionRefreshPlan(MTMVRefreshContext context, boolean canRefreshByPartitions,
List<String> partitions, String fallbackReason) {
this.context = context;
this.canRefreshByPartitions = canRefreshByPartitions;
this.partitions = partitions == null ? Lists.newArrayList() : partitions;
this.fallbackReason = fallbackReason;
}
private static PartitionRefreshPlan success(MTMVRefreshContext context, List<String> partitions) {
return new PartitionRefreshPlan(context, true, partitions, null);
}
private static PartitionRefreshPlan fallback(String fallbackReason) {
return new PartitionRefreshPlan(null, false, Lists.newArrayList(), fallbackReason);
}
}
private static class PartitionPlanningException extends Exception {
private PartitionPlanningException(String message) {
super(message);
}
private PartitionPlanningException(String message, Throwable cause) {
super(message, cause);
}
}
@SerializedName(value = "di")
private long dbId;
@SerializedName(value = "mi")
private long mtmvId;
@SerializedName("taskContext")
private MTMVTaskContext taskContext;
// What this task reports as refreshed, and what it reports as done, for the whole task rather than for
// its last attempt: an attempt takes a scope of its own and a later one would otherwise replace it,
// hiding the partitions an earlier one committed -- which the MV has already published. Sets, because
// a later attempt's scope can cover an earlier one's (a whole-MV rebuild after a per-partition one),
// and a partition counted twice would report more work than the MV has partitions. Sorted, so that the
// column a task reports with reads the same on every look. The record is an array of names either way,
// which is what it was when these were lists.
@SerializedName("needRefreshPartitions")
Set<String> needRefreshPartitions;
@SerializedName("completedPartitions")
Set<String> completedPartitions;
@SerializedName("refreshMode")
MTMVTaskRefreshMode refreshMode;
@SerializedName("lastQueryId")
String lastQueryId;
// Persisted for SHOW MTMV TASK diagnostics. It records the IVM pre-execution
// reason that caused fallback, or the hard failure reason from IVM execution.
@SerializedName("ifr")
private String ivmFallbackReason;
@SerializedName("cg")
private String computeGroup;
private MTMV mtmv;
private MTMVRelation relation;
// Written by the executing (Disruptor worker) thread via the executeCommand consumer
// callback and read by the cancel (command) thread, so it must be volatile.
private volatile StmtExecutor executor;
// What this task has committed, per MV partition: the snapshot each partition's rows were read at.
// One accumulator for the whole task rather than one per phase, because that is what the MV publishes
// at the end of it -- a phase that started from empty would publish its own work and drop the work of
// the phases before it, leaving partitions a preceding rebuild replaced looking unsynced.
private Map<String, MTMVRefreshPartitionSnapshot> partitionSnapshots = Maps.newConcurrentMap();
// The requirement each refreshed partition was read under, captured before the base tables were read
// and recorded only once that batch's data committed (see commitCapturedEpochs). In memory only: the
// journal carries the resulting states, and a replay applies those instead of recomputing anything.
private transient Map<String, Long> ivmCapturedEpochs = Maps.newConcurrentMap();
// The requirement every partition had when this refresh planned its work. What a batch records is
// clamped to it (see commitCapturedEpochs): a mark that lands after the plan must leave its partition
// dirty rather than be written back as satisfied. Empty on a path that does not plan partition work,
// which is the plain COMPLETE path -- a whole-MV rebuild replaces every partition, so whatever it read
// is what it repaired.
private transient Map<String, Long> ivmPlannedEpochs = Maps.newHashMap();
// What a partial read answered with, held back from the batch's own recording: the partitions were
// replaced from an image of a base table older than that table is, and the delta that follows is what
// brings them up to date. The incremental attempt publishes them once it has; a task whose delta does not
// run leaves them here, unrecorded. See executePartitionBasedRefresh and redeemHeldRecords.
private transient Map<String, MTMVRefreshPartitionSnapshot> snapshotsHeldUntilTheDeltaRuns = Maps.newHashMap();
private transient Map<String, Long> epochsHeldUntilTheDeltaRuns = Maps.newHashMap();
// The MV partition each captured epoch belongs to, by the name those epochs are keyed by. A retry's
// partition sync can replace a partition with another one of the same name, and what this task holds for
// that name then describes a partition that is gone; the id is how that is told apart from a partition
// this task rebuilt itself, whose rows are there and whose state is still the one it started with.
private transient Map<String, Long> capturedPartitionIds = Maps.newHashMap();
// How many partitions this refresh rebuilt because the criterion demanded it, which a strict
// INCREMENTAL request reports so that "the request was incremental but the work was not" is visible.
@SerializedName("irp")
private int ivmRebuiltPartitions;
private long mtmvSchemaChangeVersion;
// Published only after a signature-mismatch fallback succeeds and its task result is accepted.
private transient String refreshedIvmPlanSignature;
private Map<MvccTableInfo, MvccSnapshot> snapshots = Maps.newHashMap();
public MTMVTask() {
}
public MTMVTask(long dbId, long mtmvId, MTMVTaskContext taskContext) {
this.dbId = Objects.requireNonNull(dbId);
this.mtmvId = Objects.requireNonNull(mtmvId);
this.taskContext = Objects.requireNonNull(taskContext);
}
// only for test
public MTMVTask(MTMV mtmv, MTMVRelation relation, MTMVTaskContext taskContext) {
this.mtmv = mtmv;
this.relation = relation;
this.taskContext = taskContext;
}
@Override
public void run() throws JobException {
if (LOG.isDebugEnabled()) {
LOG.debug("mtmv task run, taskId: {}", super.getTaskId());
}
mtmvSchemaChangeVersion = mtmv.getSchemaChangeVersion();
ConnectContext ctx = MTMVPlanUtil.createMTMVContext(mtmv, MTMVPlanUtil.DISABLE_RULES_WHEN_RUN_MTMV_TASK);
try {
if (LOG.isDebugEnabled()) {
String taskSessionContext = ctx.getSessionVariable().toJson().toJSONString();
if (LOG.isDebugEnabled()) {
LOG.debug("mtmv task session variable, taskId: {}, session: {}",
super.getTaskId(), taskSessionContext);
}
}
// Every time a task is run, the relation is regenerated because baseTables and baseViews may change,
// such as deleting a table and creating a view with the same name
MTMVPlanUtil.QueryAnalysisResult queryAnalysis = MTMVPlanUtil.getBaseTableFromQuery(
mtmv.getQuerySql(), ctx);
this.relation = MTMVPlanUtil.generateMTMVRelation(queryAnalysis.getAllLevelTables(),
queryAnalysis.getOneLevelTables());
beforeMTMVRefresh();
List<TableIf> tableIfs = Lists.newArrayList(queryAnalysis.getAllLevelTables());
tableIfs.sort(Comparator.comparing(TableIf::getId));
// This checks whether an MV in SCHEMA_CHANGE state still matches
// its base-table schema and partition definition. It is not part of
// refresh fallback: incompatible MV definitions must fail directly.
ensureQueryUsableIfNeeded(ctx, tableIfs);
RefreshRequest request = resolveRefreshRequest();
try {
syncPartitionsIfNeeded(ctx, tableIfs);
} catch (PartitionPlanningException e) {
throw new JobException(e.getMessage(), e);
}
// Partition sync has decided which partitions exist, and nothing has read a base table yet:
// this is the point where an entry and the partition it describes become the same thing.
// Doing it any later would let a partition that sync has just added be refreshed without an
// entry, and an invalidation arriving in between would have nothing to land on.
mtmv.alignPartitionStates();
// Decided after the sync and the alignment, because the escalation it can take reads the
// partition states and only then is the partition set they describe final.
List<RefreshAttemptType> attempts = buildAttempts(request, queryAnalysis.containsOneRowRelation());
MTMVRefreshContext refreshContext = buildRefreshContext(tableIfs);
boolean disablePartitionRefresh = false;
for (RefreshAttemptType attemptType : attempts) {
switch (attemptType) {
case IVM:
AttemptResultType ivmResult = executeIvmAttempt(refreshContext, request, ctx, tableIfs);
if (ivmResult == AttemptResultType.SUCCESS) {
return;
}
if (ivmResult == AttemptResultType.FALLBACK_TO_COMPLETE) {
disablePartitionRefresh = true;
}
break;
case PARTITIONS:
if (disablePartitionRefresh) {
break;
}
if (executePartitionBasedRefresh(refreshContext, request, ctx)) {
return;
}
break;
case COMPLETE:
try {
executeCompleteAttempt(refreshContext, ctx);
} finally {
// Counted from the groups that committed, like the rebuild phase above: a
// whole-MV rebuild that failed in a later group has still replaced the groups
// before it, and those keep the epochs and snapshots they published (see
// MTMV#addTaskResult), so reporting none of them would hide the partial work
// this count exists to expose.
recordRebuiltPartitions(request);
}
return;
default:
throw new JobException("Unsupported refresh attempt type: " + attemptType);
}
}
throw new JobException("No refresh attempt succeeded for mv=" + mtmv.getName());
} catch (Throwable e) {
if (getStatus() == TaskStatus.RUNNING) {
LOG.warn("run task failed, mvName: {}, taskId: {}",
mtmv.getName(), getTaskId(), e);
throw new JobException(e.getMessage(), e);
} else {
// if status is not `RUNNING`,maybe the task was canceled, therefore, it is a normal situation
LOG.info("task [{}] interruption running, because status is [{}]", getTaskId(), getStatus());
}
} finally {
closeExecutionContext(ctx);
}
}
private void ensureQueryUsableIfNeeded(ConnectContext ctx, List<TableIf> tableIfs)
throws JobException, AnalysisException {
MetaLockUtils.readLockTables(tableIfs);
try {
if (MTMVState.SCHEMA_CHANGE.equals(mtmv.getStatus().getState())) {
MTMVPlanUtil.ensureMTMVQueryUsable(mtmv, ctx);
}
} finally {
MetaLockUtils.readUnlockTables(tableIfs);
}
}
private void syncPartitionsIfNeeded(ConnectContext ctx, List<TableIf> tableIfs)
throws JobException, AnalysisException, DdlException, PartitionPlanningException {
if (isSkipPartitionSyncDebugPointEnabled()) {
LOG.info("Skip MTMV partition synchronization for debug point, mv={}, taskId={}",
mtmv.getName(), getTaskId());
return;
}
Pair<List<String>, List<PartitionKeyDesc>> syncPartitions = null;
// lock table order by id to avoid deadlock
MetaLockUtils.readLockTables(tableIfs);
try {
if (mtmv.getMvPartitionInfo().getPartitionType() != MTMVPartitionType.SELF_MANAGE) {
Set<MTMVRelatedTableIf> pctTables = mtmv.getMvPartitionInfo().getPctTables();
for (MTMVRelatedTableIf pctTable : pctTables) {
if (!pctTable.isValidRelatedTable()) {
throw new PartitionPlanningException("MTMV " + mtmv.getName()
+ "'s pct table " + pctTable.getName()
+ " is not a valid pct table anymore, stop refreshing."
+ " e.g. Table has multiple partition columns"
+ " or including not supported transform functions.");
}
}
try {
syncPartitions = MTMVPartitionUtil.alignMvPartition(mtmv, snapshots);
} catch (Exception e) {
throw new PartitionPlanningException(e.getMessage(), e);
}
}
} finally {
MetaLockUtils.readUnlockTables(tableIfs);
}
if (syncPartitions != null) {
for (String pName : syncPartitions.first) {
MTMVPartitionUtil.dropPartition(mtmv, pName);
}
for (PartitionKeyDesc partitionKeyDesc : syncPartitions.second) {
MTMVPartitionUtil.addPartition(mtmv, partitionKeyDesc);
}
}
}
private boolean isSkipPartitionSyncDebugPointEnabled() {
String targetMvName = DebugPointUtil.getDebugParamOrDefault(
DEBUG_POINT_SKIP_PARTITION_SYNC_FILTER, "mv_name", "");
return mtmv.getName().equals(targetMvName)
&& DebugPointUtil.isEnable(DEBUG_POINT_SKIP_PARTITION_SYNC);
}
private RefreshRequest resolveRefreshRequest() throws JobException {
if (taskContext.useMvDefaultRefreshPolicy()) {
// Scheduled/on-commit/system tasks use the policy persisted on the
// MV, not the default AUTO value of a newly created task context.
RefreshMethod refreshMethod = mtmv.getRefreshInfo().getRefreshMethod();
if (refreshMethod == null) {
throw new JobException("MTMV " + mtmv.getName()
+ " has unknown refresh method, please refresh or recreate it.");
}
return new RefreshRequest(RefreshMode.valueOf(refreshMethod.name()),
mtmv.getRefreshInfo().allowFallback(), Lists.newArrayList(), false);
}
if (!CollectionUtils.isEmpty(taskContext.getPartitions())) {
// A partitionSpec is an exact manual request. It never falls back to
// COMPLETE because that would refresh more data than the user asked.
return new RefreshRequest(RefreshMode.PARTITIONS, false, taskContext.getPartitions(), true);
}
return new RefreshRequest(taskContext.getRefreshMode(), taskContext.allowFallback(),
Lists.newArrayList(), false);
}
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) {
// A PARTITIONS request that may not fall back asked for a scope, and this invalidation needs
// one it did not ask for. Refreshing the partitions it names would leave the MV in
// SCHEMA_CHANGE with rows nothing rebuilt, and widening it here would rebuild partitions the
// caller deliberately kept out -- the answer is that this request cannot do it, and the
// three that can are the ones whose scope already includes a whole-MV rebuild.
throw new JobException("The refresh baseline of MV " + mtmv.getName()
+ " was invalidated, so the partitions it needs to rebuild are not the ones a"
+ " PARTITIONS refresh names. Use COMPLETE, AUTO, or PARTITIONS FALLBACK.");
}
LOG.info("IVM MV is in SCHEMA_CHANGE, rebuilding the whole MV, mv={}, taskId={}",
mtmv.getName(), getTaskId());
return Lists.newArrayList(RefreshAttemptType.COMPLETE);
}
if (shouldUseCompleteForInitialIvmRefresh(containsOneRowRelation)) {
return Lists.newArrayList(RefreshAttemptType.COMPLETE);
}
List<RefreshAttemptType> attempts = Lists.newArrayList();
switch (request.refreshMode) {
case AUTO:
// ALTER excluded_trigger_tables clears a successful baseline without changing the MV state.
if (!mtmv.isIvm() && !mtmv.hasRefreshSnapshot()
&& mtmv.getStatus().getState() == MTMVState.NORMAL
&& mtmv.getStatus().getRefreshState() == MTMVRefreshState.SUCCESS) {
attempts.add(RefreshAttemptType.COMPLETE);
break;
}
if (mtmv.isIvm()) {
attempts.add(RefreshAttemptType.IVM);
}
// AUTO always ends with COMPLETE. For an MV defined REFRESH COMPLETE,
// skip the PARTITIONS attempt: its sync check treats base tables that are
// not MTMVRelatedTableIf (external tables, views) as always synchronous,
// so the refresh would be skipped forever after the first build.
if (mtmv.getRefreshInfo().getRefreshMethod() != RefreshMethod.COMPLETE) {
attempts.add(RefreshAttemptType.PARTITIONS);
} else {
LOG.info("AUTO refresh of COMPLETE-method mv={} skips partition-sync check "
+ "and refreshes all directly, taskId={}", mtmv.getName(), super.getTaskId());
}
attempts.add(RefreshAttemptType.COMPLETE);
break;
case INCREMENTAL:
attempts.add(RefreshAttemptType.IVM);
if (request.allowFallback) {
attempts.add(RefreshAttemptType.PARTITIONS);
attempts.add(RefreshAttemptType.COMPLETE);
}
break;
case PARTITIONS:
attempts.add(RefreshAttemptType.PARTITIONS);
if (!request.explicitPartitions && request.allowFallback) {
attempts.add(RefreshAttemptType.COMPLETE);
}
break;
case COMPLETE:
attempts.add(RefreshAttemptType.COMPLETE);
break;
default:
throw new IllegalStateException("Unsupported refresh mode: " + request.refreshMode);
}
// A base table the plan scans without a usable stream makes the incremental attempt fail: its
// rewrite reads the stream of every table it scans. Only COMPLETE reconciles streams, so a request
// that would try the incremental path first goes there directly, and the attempt that cannot
// succeed -- with the baseline barrier it writes before it starts -- stays out of the way.
//
// Only that attempt is judged here. A partition refresh reads a narrower set, and how narrow
// depends on the partitions it plans: a PCT table that no refreshed partition's mapping names
// keeps its place in the plan as an ordinary scan, and an excluded trigger table is never read
// through a stream. It checks its own scope once it has planned (see
// hasUnusableIvmStreamForPartitions), so judging it here by the whole plan would send a partition
// refresh that would have worked to COMPLETE.
if (request.allowFallback && mtmv.isIvm() && attempts.contains(RefreshAttemptType.IVM)
&& hasUnusableIvmStream()) {
ivmFallbackReason = IvmFailureReason.STREAM_UNSUPPORTED.name();
LOG.warn("IVM stream is unusable, mv={}, taskId={}. Continuing with COMPLETE refresh.",
mtmv.getName(), getTaskId());
return Lists.newArrayList(RefreshAttemptType.COMPLETE);
}
// Every partition either needs a rebuild or was never filled, and at least one needs a rebuild:
// COMPLETE then does nothing the per-partition routing would not, in one read of the MV.
//
// Only for a request that may fall back, like the stream shortcut above. The routing those
// partitions would take is the incremental attempt, which rebuilds every one of them and then finds
// no scope left to catch up -- and it reads the streams, so a request that may not fall back still
// fails there when one is unusable. Answering COMPLETE instead does what such a request forbids: it
// reconciles the streams and resets their baselines, which is the same reset the shortcut above
// gates on allowing a fallback. A strict request reaches the incremental attempt, as its intent is
// stated there.
//
// No clause for an explicit partition list: the attempt list says it already. Such a list is
// rejected for an IVM MV when the statement is analyzed, and it becomes a PARTITIONS request here,
// which never reaches the incremental attempt -- so an attempt list that holds one cannot hold the
// other. The schema-change shortcut above needs its own clause because it is judged for requests
// that do name partitions.
if (request.allowFallback && attempts.contains(RefreshAttemptType.IVM)
&& mtmv.allPartitionsNeedRebuild()) {
LOG.info("Every MV partition needs a rebuild or has no data yet, mv={}, taskId={}. "
+ "Continuing with COMPLETE refresh.", mtmv.getName(), getTaskId());
return Lists.newArrayList(RefreshAttemptType.COMPLETE);
}
return attempts;
}
/**
* Notes that this refresh rebuilds partitions the request did not ask to rebuild, which is what a
* strict INCREMENTAL request cannot tell from its result otherwise: it reports the count, and a request
* that asked for a complete refresh reports nothing because rebuilding everything is what it asked for.
*
* <p>What it counts is the partitions this task has replaced by now -- the batches that committed, read
* from the same accumulator the task reports its progress with, so a refresh that failed part-way
* through cannot report partitions it never replaced. None of them is the empty case: an attempt that
* found nothing to refresh, and one the stream reconciliation threw out of before it reached the
* executor, both report on a task that has committed nothing. That has no count, and it reads as zero
* rather than throwing where the refresh is being reported.
*
* <p>The value describes the refresh rather than its last attempt: a later attempt's count covers the
* partitions an earlier one rebuilt -- a whole-MV attempt replaces every partition it commits -- so the
* largest any of them reported is the number this refresh replaced, and an attempt that replaced fewer
* of them must not lower it.
*/
private void recordRebuiltPartitions(RefreshRequest request) {
// Only an IVM MV has a baseline to rebuild: a plain MV's COMPLETE is the only way it refreshes at
// all, so reporting it there would put a rebuild count on every ordinary refresh.
if (!mtmv.isIvm() || request.refreshMode == RefreshMode.COMPLETE) {
return;
}
int replaced = CollectionUtils.isEmpty(completedPartitions) ? 0 : completedPartitions.size();
ivmRebuiltPartitions = Math.max(ivmRebuiltPartitions, replaced);
}
private boolean shouldUseCompleteForInitialIvmRefresh(boolean containsOneRowRelation) {
if (!mtmv.isIvm() || mtmv.hasRefreshSnapshot()) {
return false;
}
// Excluded trigger tables produce no delta stream in the incremental rewrite, so an
// INCREMENTAL first refresh would never read their pre-existing rows. Whenever the
// MV has excluded sources, the first refresh must build the full COMPLETE baseline
// first (covering manual INCREMENTAL as well); afterwards the snapshot exists and
// incremental refreshes only track the remaining streamed tables.
return taskContext.getTriggerMode() != MTMVTaskTriggerMode.MANUAL
|| !CollectionUtils.isEmpty(mtmv.getExcludedTriggerTables())
|| containsOneRowRelation;
}
private PartitionRefreshPlan planPartitionRefresh(MTMVRefreshContext context,
RefreshRequest request) throws AnalysisException {
if (request.explicitPartitions) {
return PartitionRefreshPlan.success(context, request.partitions);
}
// The partitions the criterion says have to be rebuilt, whatever the snapshots say: a snapshot
// records what a refresh last read, and a raised requirement says that read cannot be caught up.
// The two disagree in the one state a whole-MV rebuild that did not finish leaves behind -- every
// requirement raised, every snapshot kept -- and planning by snapshots alone would report
// NOT_REFRESH over partitions whose rows nothing repaired. A partition refresh is what rebuilds
// them (the incremental path rebuilds its own dirty set the same way), so they are planned here
// rather than left to an attempt this request may never take.
Set<String> rebuildRequired = mtmv.getPartitionsNeedingRebuild();
boolean fresh;
try {
fresh = MTMVPartitionUtil.isMTMVSync(context, relation.getBaseTablesOneLevelAndFromView(),
mtmv.getExcludedTriggerTables());
} catch (Exception e) {
return PartitionRefreshPlan.fallback(e.getMessage());
}
if (fresh && rebuildRequired.isEmpty()) {
return PartitionRefreshPlan.success(context, Lists.newArrayList());
}
if (mtmv.getMvPartitionInfo().getPartitionType() == MTMVPartitionType.SELF_MANAGE) {
// Keep this inside the PARTITIONS attempt so PARTITIONS FALLBACK and
// AUTO can still continue to COMPLETE for non-partitioned MVs.
return PartitionRefreshPlan.fallback(
"The partition method of this asynchronous materialized view "
+ "does not support refreshing by partition");
}
try {
Set<String> planned = Sets.newLinkedHashSet(MTMVPartitionUtil.getMTMVNeedRefreshPartitions(context,
relation.getBaseTablesOneLevelAndFromView()));
planned.addAll(rebuildRequired);
return PartitionRefreshPlan.success(context,
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>A plan that covers the whole MV is left alone, because for one of those the second read is not the
* same read. What a whole-MV plan gets from the executor is a RESET read of the tables the MV does not
* partition by: it reads them as they are now, and it advances their stream offsets to the end, which
* only a read that saw the whole table may do. Taking partitions out of that plan makes it a partial
* one, and a partial read of those tables is a snapshot of the offset instead: it reads them as they
* were when the offset was recorded. A partition this task replaced through such a read and is then
* skipped is left holding that old image while the refresh replaces the rest with the current one, has
* moved the offset past the difference, and records the current state for every partition it planned --
* which is a partition that is wrong and looks up to date, and nothing left to repair it.
*
* <p>Leaving the whole plan alone costs the second read in the case where the earlier phase had itself
* planned the whole MV. That phase read through RESET, so the partitions it replaced are current and
* skipping them would have been safe; telling the two apart would take remembering, per partition, how
* its replacement read, and the price of not remembering is one overwrite.
*
* <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()
|| Sets.newHashSet(plannedPartitions).equals(mtmv.getPartitionNames())) {
return plannedPartitions;
}
List<String> remaining = Lists.newArrayListWithCapacity(plannedPartitions.size());
for (String partitionName : plannedPartitions) {
if (!partitionSnapshots.containsKey(partitionName)) {
remaining.add(partitionName);
}
}
return remaining;
}
private MTMVRefreshContext buildRefreshContext(List<TableIf> tableIfs) throws AnalysisException {
MetaLockUtils.readLockTables(tableIfs);
try {
return MTMVRefreshContext.buildContext(mtmv, Maps.newHashMap(), snapshots);
} finally {
MetaLockUtils.readUnlockTables(tableIfs);
}
}
private void executeCompleteAttempt(MTMVRefreshContext context, ConnectContext ctx)
throws JobException, AnalysisException {
List<String> partitions = Lists.newArrayList(mtmv.getPartitionNames());
// Reported here rather than by the executor alone: a whole-MV attempt covers every partition from
// this point on, and the reconciliation below can throw before the executor is reached -- leaving a
// report that names what the attempt was about to rebuild rather than one that names nothing.
recordRefreshScope(partitions);
// Nothing to clamp a captured epoch against: a whole-MV rebuild replaces every partition, so what
// this refresh read is what it repaired. Dropped rather than kept, so the partitions an earlier
// partial rebuild did replace do not stay looking like they still owe one.
this.ivmPlannedEpochs = Maps.newHashMap();
this.refreshMode = generateRefreshMode(partitions);
if (refreshMode == MTMVTaskRefreshMode.NOT_REFRESH) {
return;
}
// Marked before the streams are reconciled, not merely before the rebuild: reconciling durably
// creates a replacement stream whose historical rows are read as an append, so a crash after that
// create and before this rebuild publishes its epochs would leave a populated MV with clean epochs
// and a usable stream -- and the next incremental refresh would add those rows to the old baseline
// again. The requirement is raised as a real mark on the partitions, which is what "these must be
// rebuilt" means to the refresh, and it is journaled as the whole map before the reconcile below.
mtmv.markPartitionsForRebuild(Sets.newHashSet(partitions));
// A complete rebuild resets the stream baselines, so reconcile missing or unusable streams
// before reading anything. Only COMPLETE may do this: a stream baseline is global, resetting
// it during a partial refresh would corrupt the partitions that refresh does not touch.
if (mtmv.isIvm()) {
reconcileIvmStreams(ctx);
}
executePartitionBasedRefresh(context, RefreshMode.COMPLETE, ctx, partitions);
}
private AttemptResultType executeIvmAttempt(MTMVRefreshContext refreshContext,
RefreshRequest request, ConnectContext ctx, List<TableIf> tableIfs)
throws JobException, AnalysisException {
if (!mtmv.isIvm()) {
throw new JobException("Cannot use " + request.refreshMode
+ " refresh on a materialized view without INCREMENTAL capability.");
}
// A strict INCREMENTAL request must reach IVM so the stream can establish its initial baseline.
// Only fallback-enabled requests may rebuild a missing or invalidated snapshot with COMPLETE.
if (!mtmv.hasRefreshSnapshot() && request.allowFallback) {
ivmFallbackReason = "INCOMPLETE_REFRESH_SNAPSHOT";
LOG.warn("IVM refresh fell back for mv={}, reason=INCOMPLETE_REFRESH_SNAPSHOT, taskId={}. "
+ "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, 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();
return AttemptResultType.SUCCESS;
}
if (ivmResult.getFailureReason() != IvmFailureReason.MV_PARTITION_NOT_FOUND) {
return handleIvmFallbackResult(ivmResult, request);
}
if (partitionSyncRetryCount + 1 >= ivmAttemptLimit) {
break;
}
try {
syncPartitionsIfNeeded(ctx, tableIfs);
// The retry can add a partition that did not exist at the first alignment. It has to get
// its entry before the retried refresh reads a base table, or an invalidation arriving
// in between would have nothing to land on for rows this task is about to write.
mtmv.alignPartitionStates();
// Alignment gives a partition it creates {0, 1} -- behind its requirement -- and the routing
// decision was taken before that partition existed. Without a fresh read the retried attempt
// would hand it to the delta path with no ceiling to be clamped against, and its capture,
// which is all that path can produce, would be recorded as the partition being caught up:
// a baseline it never received, and no repairing delta to notice.
adoptPartitionsCreatedByTheRetry(dirtyPartitions);
currentRefreshContext = buildRefreshContext(tableIfs);
} catch (Exception e) {
throw new JobException("Failed to synchronize MV partitions before IVM retry for mv="
+ mtmv.getName(), e);
}
LOG.warn("Retrying IVM refresh after synchronizing MV partitions, mv={}, attempt={}/{}, taskId={}",
mtmv.getName(), partitionSyncRetryCount + 1, ivmAttemptLimit, getTaskId());
}
throw new JobException("IVM refresh could not recover missing MV partition for mv="
+ mtmv.getName() + ", detail=" + ivmResult.getDetailMessage());
}
private IvmIncrRefreshResult executeSingleIvmAttempt(MTMVRefreshContext refreshContext,
Set<String> dirtyPartitions)
throws JobException {
// Determine which partitions need refresh, same as partition-based flow. The partitions the
// rebuild above handled are taken out: an incremental refresh of one of them would record it as
// caught up while its rows are exactly what the rebuild had to replace. This is what the attempt
// may record, which is not the same question as what it has to run for; see
// hasRecordsHeldForTheDelta.
Set<String> incrementalScope = Sets.newLinkedHashSet(MTMVPartitionUtil.getMTMVNeedRefreshPartitions(
refreshContext, relation.getBaseTablesOneLevelAndFromView()));
incrementalScope.removeAll(dirtyPartitions);
recordRefreshScope(incrementalScope);
if (CollectionUtils.isEmpty(incrementalScope) && !hasRecordsHeldForTheDelta()) {
LOG.info("IVM incremental refresh skipped for mv={}: all partitions are synced, taskId={}",
mtmv.getName(), getTaskId());
return IvmIncrRefreshResult.success();
}
IvmIncrRefreshManager ivmIncrRefreshManager = new IvmIncrRefreshManager();
// Capture base table snapshots under read lock before execution, same as
// partition-based refresh. This ensures snapshot versions are consistent
// with the data the INSERT will read.
Map<String, MTMVRefreshPartitionSnapshot> capturedSnapshots;
try {
capturedSnapshots = MTMVPartitionUtil.generatePartitionSnapshots(
refreshContext, relation.getBaseTablesOneLevelAndFromView(),
Sets.newHashSet(incrementalScope));
} catch (Exception e) {
throw new JobException("IVM snapshot generation failed for mv=" + mtmv.getName(), e);
}
// The requirement these partitions are read under, captured before the read inside doRefresh and
// recorded only if the refresh commits; see captureLatestEpochs.
Map<String, Long> capturedEpochs = captureLatestEpochs(incrementalScope);
IvmIncrRefreshResult ivmResult;
try {
ivmResult = executeWithRetry(() -> {
ConnectContext ivmConnectContext = MTMVPlanUtil.createMTMVContext(mtmv,
MTMVPlanUtil.DISABLE_RULES_WHEN_RUN_MTMV_TASK);
try {
setupComputeGroup(ivmConnectContext);
IvmIncrRefreshContext ivmIncrRefreshContext = new IvmIncrRefreshContext(mtmv,
ivmConnectContext,
getRefreshAuditStmt(RefreshMode.INCREMENTAL, Sets.newHashSet(incrementalScope)),
this::recordQueryId,
this::registerExecutor);
mtmv.validateIvmRefreshStart(mtmvSchemaChangeVersion);
return ivmIncrRefreshManager.doRefresh(ivmIncrRefreshContext);
} finally {
closeExecutionContext(ivmConnectContext);
}
}, "IVM refresh");
} catch (Exception e) {
throw new JobException("IVM incremental refresh failed for mv=" + mtmv.getName()
+ ", detail=" + Util.getRootCauseMessage(e), e);
}
if (ivmResult.isSuccess()) {
this.partitionSnapshots.putAll(capturedSnapshots);
recordRefreshCompleted(incrementalScope);
commitCapturedEpochs(capturedEpochs);
LOG.info("IVM incremental refresh succeeded for mv={}, taskId={}",
mtmv.getName(), getTaskId());
}
return ivmResult;
}
/**
* Whether the delta is what a partition this task replaced is still waiting for: the records the rebuild
* held back, whose publication is the delta's to earn.
*
* <p>An attempt whose scope comes out empty is one where every partition that needs a refresh is one the
* rebuild above replaced -- they are taken out of the scope on purpose. Skipping the delta then would
* skip it for exactly the partitions it is the only repair of: what the rebuild read of a table the MV
* does not partition by is the image as of the stream offset, and the delta is what brings that table up
* to date for them. Their records stay unrecorded until it has run, so the attempt runs for them even
* with nothing to record -- the scope is what it records, and an empty one records nothing.
*
* <p>The delta does not need the scope to do that work: 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.
*
* <p>One of the two held maps answers for both: a batch holds its epochs and its snapshots together,
* and redeeming publishes the pair -- which is also what the two are read as.
*/
private boolean hasRecordsHeldForTheDelta() {
return !epochsHeldUntilTheDeltaRuns.isEmpty();
}
/**
* Adds a phase's scope to what this task reports as refreshed, and the partitions it committed to what
* this task reports as done. Both are the task's; see the fields for why.
*/
private void recordRefreshScope(Collection<String> partitions) {
// Created by the first phase that records one: a task that has not refreshed anything reports
// nothing, which is the same state as one that has not run yet. Concurrent because the columns it
// feeds are read while the worker fills them -- the tasks() table function reports a running task
// -- and ordered because what this reports is persisted and has to read the same on every look.
if (needRefreshPartitions == null) {
needRefreshPartitions = new ConcurrentSkipListSet<>();
}
needRefreshPartitions.addAll(partitions);
}
private void recordRefreshCompleted(Collection<String> partitions) {
if (completedPartitions == null) {
completedPartitions = new ConcurrentSkipListSet<>();
}
completedPartitions.addAll(partitions);
}
/**
* Says durably that the parts of the scope that do not name a rebuild requirement have to be rebuilt, and
* brings the epochs this phase records in line with it.
*
* <p>Raised here rather than by each caller of the executor, like the scope above: this is the phase that
* replaces partitions, and a caller that forgot would leave a partition whose rows were never published
* looking caught up. The partitions a caller has already made dirty are left as they are, which is what
* makes raising it here harmless for them: a whole-MV attempt marks its scope before it reconciles the
* streams, and the incremental attempt rebuilds the partitions an invalidation marked.
* See MTMV#raiseRebuildRequirement.
*
* <p>What that call reports is what a partition it raised now names, and this phase is clamped to it. The
* clamp cannot stay at what an earlier attempt planned: that value sits below the requirement this phase
* has just raised, so the epochs recorded here would leave the partition dirty after it was replaced, and
* every refresh after it would rebuild the same partitions again.
*
* <p>A partition that already named a requirement keeps the entry it has, and is deliberately not moved
* up to what it names now. The entry is the value the routing decision saw, and it is what keeps a mark
* landing between that decision and this phase's read from being recorded as met by a replacement that
* read before the change it made. Leaving it where it is costs one rebuild; moving it up could cost the
* change.
*/
private void raiseRequirementForRefreshScope(Collection<String> partitions) {
if (!mtmv.isIvm()) {
// A plain MV has no streams to read and its epochs record nothing: the sync criterion plans it
// again on its own until its snapshots are published.
return;
}
ivmPlannedEpochs.putAll(mtmv.raiseRebuildRequirement(Sets.newHashSet(partitions)));
}
/**
* Brings the routing decision up to date after a retry has synchronized and aligned the MV's partitions.
*
* <p>A partition the alignment creates is dirty by construction -- {@code {0, 1}}, behind its
* requirement -- and the decision was taken before it existed. Reading the states again is what makes
* the retried attempt treat it as such: it joins the dirty set, so the incremental attempt leaves it
* out rather than recording what a delta captured as the partition being caught up, and it gets the
* entry the batches are clamped against, so a mark landing later in this task cannot be written back as
* satisfied either.
*
* <p>A partition this leaves dirty is not rebuilt here. The rebuild phase has already run, and the
* partition the alignment created holds no rows yet -- there is nothing to replace in it -- so what it
* needs is a build, which is what the next refresh's rebuild gives it.
*
* <p>The planned value of a partition that already had one is kept: that is the value the routing
* decision was made on, which is what the clamp is for.
*/
private void adoptPartitionsCreatedByTheRetry(Set<String> dirtyPartitions) {
Set<String> livePartitionNames = mtmv.getPartitionNames();
for (Entry<String, MTMVPartitionState> entry : mtmv.getPartitionStates().entrySet()) {
if (!livePartitionNames.contains(entry.getKey())) {
continue;
}
ivmPlannedEpochs.putIfAbsent(entry.getKey(), entry.getValue().getLatestEpoch());
if (entry.getValue().isDirty()) {
dirtyPartitions.add(entry.getKey());
}
// What this task holds for a name is that name's only while the partition behind it is the same
// one: the retry's partition sync can drop a partition and add another of the same name back, and
// then the captures and snapshots describe a partition that is gone. Writing them back would
// credit the new one with what the old one held, which is worse than a wrong number -- a partition
// clean at an epoch a later change only raises to is one no refresh rebuilds, so the rows the
// recreation removed would be published as current. Records this task held back are the same
// question: nothing published them yet, and a delta that succeeds after this point is what would
// redeem them onto the new partition, which has a baseline it never received.
//
// Told apart by id rather than by the state: a partition this task rebuilt has its rows and its
// capture, and its state still reads as the one an entry starts with until the task result writes
// the epochs back, so the state cannot say whether the name means the same partition.
// A partition the sync dropped has no state to walk here, so this one is live; it could still be
// dropped by a concurrent DDL, which is the one case where there is nothing to compare with.
Partition partition = mtmv.getPartition(entry.getKey());
Long capturedId = capturedPartitionIds.get(entry.getKey());
if (partition != null && capturedId != null && capturedId.longValue() != partition.getId()) {
ivmCapturedEpochs.remove(entry.getKey());
partitionSnapshots.remove(entry.getKey());
epochsHeldUntilTheDeltaRuns.remove(entry.getKey());
snapshotsHeldUntilTheDeltaRuns.remove(entry.getKey());
capturedPartitionIds.remove(entry.getKey());
}
}
}
/**
* Captures the requirement these partitions are about to be read under: the epoch in force at the
* moment the refresh starts reading, which is what the data it writes will be described by.
*
* <p>A refresh writes back the requirement it captured, not the one in force when it finishes, which
* is what keeps an invalidation that arrives while the refresh runs from being swallowed: the
* requirement it raises stays above the value the task writes, so the partition still counts as
* needing a rebuild.
*
* <p>It has to run before the base tables are read and never after. An epoch captured after the read
* could claim data newer than what the read saw, and the partition would then look caught up while
* it holds rows from before the change.
*
* <p>The caller keeps the result and hands it to {@link #commitCapturedEpochs} only once that batch's
* data has committed. Recording it here would credit a batch whose write never happened with data
* that does not exist, which is the one direction the epoch must never be wrong in.
*
* <p>A non-IVM MV carries no states, so this captures nothing for it.
*/
private Map<String, Long> captureLatestEpochs(Set<String> partitionNames) {
if (CollectionUtils.isEmpty(partitionNames)) {
return Maps.newHashMap();
}
return mtmv.getLatestEpochs(partitionNames);
}
/**
* Commits the captured epochs of a batch whose data has landed, so its work is not repeated after a
* restart; see {@link #captureLatestEpochs} for what a captured value is. A partition read by two
* phases of one task keeps the higher value, which is the requirement its surviving data was read
* under.
*/
private void commitCapturedEpochs(Map<String, Long> capturedEpochs) {
for (Entry<String, Long> entry : capturedEpochs.entrySet()) {
ivmCapturedEpochs.merge(entry.getKey(), plannedCeiling(entry), Math::max);
recordCapturedPartitionId(entry.getKey());
}
}
/**
* Notes the partition the record for this name was taken under, so a later phase can tell whether the name
* still means that partition; see {@link #adoptPartitionsCreatedByTheRetry}.
*
* <p>Recorded for records this task holds back as well as for the ones it has committed -- both are its,
* and a retry that replaces the partition behind a name has to be able to drop either. A held record
* without an id here would be redeemed onto the partition that took the name, crediting it with a
* baseline it never received.
*/
private void recordCapturedPartitionId(String partitionName) {
Partition partition = mtmv.getPartition(partitionName);
if (partition != null) {
capturedPartitionIds.put(partitionName, partition.getId());
}
}
/**
* Holds a batch's records back until the delta that brings the tables it read up to date has run, and
* notes the partitions they were taken under -- see {@link #redeemHeldRecords} for what publishing them
* waits for, and {@link #recordCapturedPartitionId} for why the ids are recorded here too.
*
* <p>Named by both maps: a batch holds a snapshot for every partition it replaced, and an epoch only for
* the ones that already had an entry to capture, so the names are the union of the two rather than the
* epochs alone.
*/
private void holdCapturedEpochs(Map<String, MTMVRefreshPartitionSnapshot> snapshots,
Map<String, Long> capturedEpochs) {
snapshotsHeldUntilTheDeltaRuns.putAll(snapshots);
epochsHeldUntilTheDeltaRuns.putAll(capturedEpochs);
for (String partitionName : Sets.union(snapshots.keySet(), capturedEpochs.keySet())) {
recordCapturedPartitionId(partitionName);
}
}
/**
* Publishes what a partial read held back, once the delta that brings those tables up to date has run.
*
* <p>Before that the records are not this task's to publish: the partitions hold the image of a base table
* older than the table is, and recording them would say they are current while they are not. After it they
* are, and holding them any longer is not free -- a partition left unrecorded is one the next refresh
* rebuilds in full, every time a table it reads keeps changing.
*/
private void redeemHeldRecords() {
if (epochsHeldUntilTheDeltaRuns.isEmpty()) {
return;
}
partitionSnapshots.putAll(snapshotsHeldUntilTheDeltaRuns);
commitCapturedEpochs(epochsHeldUntilTheDeltaRuns);
epochsHeldUntilTheDeltaRuns = Maps.newHashMap();
snapshotsHeldUntilTheDeltaRuns = Maps.newHashMap();
}
/**
* The epoch to record for a captured partition: the one it was read at, or the one the routing decision
* saw when that is lower -- see {@link #captureLatestEpochs} for the rule this serves.
*
* <p>The clamp is what keeps an invalidation that lands between that decision and this read from being
* recorded as satisfied: the delta cannot remove the rows it made unusable, so the partition has to
* stay dirty for the next refresh to rebuild.
*/
private long plannedCeiling(Entry<String, Long> captured) {
Long planned = ivmPlannedEpochs.get(captured.getKey());
return planned == null ? captured.getValue() : Math.min(captured.getValue(), planned);
}
private AttemptResultType handleIvmFallbackResult(IvmIncrRefreshResult ivmResult, RefreshRequest request)
throws JobException {
ivmFallbackReason = ivmResult.getFailureReason().name();
if (!request.allowFallback) {
throw new JobException(
"IVM incremental refresh failed for mv=" + mtmv.getName()
+ ", reason=" + ivmResult.getFailureReason()
+ ", detail=" + ivmResult.getDetailMessage());
}
if (ivmResult.getFailureReason() == IvmFailureReason.BINLOG_BROKEN) {
// The previous task already entered the IVM execution phase. If
// fallback is allowed, jump directly to COMPLETE recovery instead of
// trying PARTITIONS first.
LOG.warn("IVM previous run incomplete for mv={}, taskId={}. Continuing with COMPLETE recovery.",
mtmv.getName(), getTaskId());
return AttemptResultType.FALLBACK_TO_COMPLETE;
}
if (ivmResult.getFailureReason().requiresCompleteRefresh()) {
LOG.warn("IVM refresh fell back for mv={}, reason={}, detail={}, taskId={}. "
+ "Continuing with COMPLETE refresh.",
mtmv.getName(), ivmResult.getFailureReason(),
ivmResult.getDetailMessage(), getTaskId());
return AttemptResultType.FALLBACK_TO_COMPLETE;
}
LOG.warn("IVM refresh fell back for mv={}, reason={}, detail={}, taskId={}. "
+ "Continuing with partition-based refresh.",
mtmv.getName(), ivmResult.getFailureReason(),
ivmResult.getDetailMessage(), getTaskId());
return AttemptResultType.FALLBACK_ALLOWED;
}
private boolean executePartitionBasedRefresh(MTMVRefreshContext refreshContext,
RefreshRequest request, ConnectContext ctx) throws JobException, AnalysisException {
PartitionRefreshPlan partitionPlan = planPartitionRefresh(refreshContext, request);
if (!partitionPlan.canRefreshByPartitions) {
if (request.allowFallback) {
LOG.warn("MTMV partition refresh fell back for mv={}, reason={}, taskId={}",
mtmv.getName(), partitionPlan.fallbackReason, getTaskId());
return false;
}
throw new JobException(partitionPlan.fallbackReason);
}
recordRefreshScope(partitionPlan.partitions);
this.refreshMode = generateRefreshMode(partitionPlan.partitions);
// This attempt now knows which partitions it refreshes, and with them which streams it reads.
// Judged here rather than in buildAttempts because only the plan knows that scope, and judged
// before the NOT_REFRESH return below so that a fallback-capable request still reaches the only
// attempt that reconciles streams. Falling back here continues to COMPLETE, which is what repairs
// them; a request that may not fall back fails instead of quietly refreshing less than it asked.
if (mtmv.isIvm()
&& hasUnusableIvmStreamForPartitions(partitionPlan.context, partitionPlan.partitions)) {
if (!request.allowFallback) {
throw new JobException("IVM stream is unusable for the partitions of this refresh, mv="
+ mtmv.getName());
}
ivmFallbackReason = IvmFailureReason.STREAM_UNSUPPORTED.name();
LOG.warn("IVM stream is unusable for the partitions this refresh plans, mv={}, taskId={}. "
+ "Continuing with COMPLETE refresh.", mtmv.getName(), getTaskId());
return false;
}
if (refreshMode == MTMVTaskRefreshMode.NOT_REFRESH) {
return true;
}
executePartitionBasedRefresh(partitionPlan.context, RefreshMode.PARTITIONS, ctx,
partitionPlan.partitions);
return true;
}
private void executePartitionBasedRefresh(MTMVRefreshContext context, RefreshMode refreshMode,
ConnectContext ctx, List<String> partitions)
throws JobException, AnalysisException {
// Reported here rather than by the callers alone: this is the phase that runs, so this is the scope
// the report cannot do without -- a caller that forgot would leave the partitions it committed in
// the completed side and nothing in the scope, which reads as a refresh of no partitions. The
// callers name a scope as well, and only where they decide one and may not get here: a whole-MV
// attempt before it reconciles the streams, and a partition plan before it judges their streams.
recordRefreshScope(partitions);
// Whatever an earlier phase held back is this phase's to replace: it is about to read the same tables.
epochsHeldUntilTheDeltaRuns = Maps.newHashMap();
snapshotsHeldUntilTheDeltaRuns = Maps.newHashMap();
// The durable half of what this phase is about to do, raised before it reads anything. An overwrite
// commits the rows into temporary partitions and publishes them with a swap afterwards, so a refresh
// that dies between the two leaves the live partitions holding the rows they had while the streams
// its read consumed are already advanced -- and the epochs it would have written ride with a result
// that never came. See raiseRequirementForRefreshScope.
raiseRequirementForRefreshScope(partitions);
boolean useIvmFallbackStreams = mtmv.isIvm();
Map<TableIf, String> tableWithPartKey = getIncrementalTableMap();
try {
// Snapshot persistence happens after refresh partitions are split into execution groups. Load the
// complete union here so the default one-partition group size cannot turn a large Hive MTMV into
// one metadata request per MV partition; generatePartitionSnapshots reuses this context cache.
context.preparePartitionSnapshots(Sets.newHashSet(partitions));
} catch (Exception e) {
// Preloading is only a batching optimization. Retrying through the existing per-group load below
// preserves completed-group progress when a later chunk of the union fails.
LOG.warn("Failed to preload partition snapshots for mv={}, taskId={}; "
+ "falling back to per-group loading", mtmv.getName(), getTaskId(), e);
}
int refreshPartitionNum = mtmv.getRefreshPartitionNum();
long execNum = (partitions.size() / refreshPartitionNum) + ((partitions.size()
% refreshPartitionNum) > 0 ? 1 : 0);
boolean refreshAllPartitions = Sets.newHashSet(partitions).equals(mtmv.getPartitionNames());
// Only a COMPLETE that the signature-mismatch fallback asked for publishes its signature. Widening
// this to every COMPLETE an IVM MV runs -- the escalation an invalidated MV takes is the other one
// -- is not the free saving it looks like: the fresh signature lets the following refresh take the
// incremental path instead of that fallback, and the incremental then re-applies rows the COMPLETE
// had just rebuilt. Three suites report their delta twice once it is widened (test_ivm_snapshot,
// test_ivm_bitmap_agg_2, test_ivm_agg_array_1), so the second COMPLETE is protection rather than
// waste. What has to be understood first is how the incremental path baselines itself against a
// plan that has changed; until then this condition stays narrow. Non-IVM MVs produce no signature.
boolean capturePlanSignature = refreshMode == RefreshMode.COMPLETE
&& IvmFailureReason.PLAN_SIGNATURE_MISMATCH.name().equals(ivmFallbackReason);
IvmPlanSignature refreshedPlanSignature = null;
for (int i = 0; i < execNum; i++) {
int start = i * refreshPartitionNum;
int end = start + refreshPartitionNum;
Set<String> execPartitionNames = Sets.newHashSet(partitions
.subList(start, Math.min(end, partitions.size())));
Map<BaseTableInfo, Set<Long>> batchResetPartitionIds = useIvmFallbackStreams
? collectPctResetPartitionIds(context, execPartitionNames) : Maps.newHashMap();
Optional<IvmRewriteContext> rewriteContext = Optional.empty();
if (useIvmFallbackStreams) {
StreamReadMode nonPctReadMode = i == 0 && refreshAllPartitions
? StreamReadMode.RESET : StreamReadMode.SNAPSHOT;
rewriteContext = Optional.of(
IvmRewriteContext.full(mtmv, batchResetPartitionIds, nonPctReadMode));
}
// The requirement this batch is read under, captured before the read below and recorded once
// the read's data has committed, next to its snapshots. Capturing it per batch keeps an
// invalidation that arrives during the refresh from holding back the whole round: only the
// batches already read keep a requirement above their captured value.
Map<String, Long> batchCapturedEpochs = captureLatestEpochs(execPartitionNames);
// need get names before exec
Map<String, MTMVRefreshPartitionSnapshot> execPartitionSnapshots = MTMVPartitionUtil
.generatePartitionSnapshots(context, relation.getBaseTablesOneLevelAndFromView(),
execPartitionNames);
try {
IvmPlanSignature batchPlanSignature = refreshPartitionsWithRetry(
execPartitionNames, tableWithPartKey, rewriteContext, refreshMode);
if (capturePlanSignature) {
batchPlanSignature = Objects.requireNonNull(batchPlanSignature,
"IVM COMPLETE refresh did not produce a plan signature");
if (refreshedPlanSignature == null) {
refreshedPlanSignature = batchPlanSignature;
} else if (!refreshedPlanSignature.getSha256().equals(batchPlanSignature.getSha256())) {
throw new JobException("IVM COMPLETE refresh generated inconsistent plan signatures, mv="
+ mtmv.getName());
}
}
} catch (Exception e) {
LOG.error("Execution failed after retries, mvName: {}, taskId: {}",
mtmv.getName(), getTaskId(), e);
throw new JobException(e.getMessage(), e);
}
recordRefreshCompleted(execPartitionNames);
// What this batch may record: everything it replaced, unless the read behind it answered with a
// base table the MV does not partition by as of an older state than that table is in now -- which
// a partial read does for the partitions holding data the stream offset has not consumed. For
// those it is the delta that follows that brings the table up to date, and it does reach them:
// the partitions it writes are the ones the change reaches, and they are among them. What keeps
// them out of the delta's scope is this attempt's recording, not the delta's work; see
// hasRecordsHeldForTheDelta. Recording them here would say they hold the table's current state
// when they hold that image, and nothing would plan them again. Left unrecorded, the requirement
// stays, the MV keeps the snapshot it has rather than one claiming the current state, and a later
// refresh rebuilds them -- by then the offset has been consumed, so a rebuild reads the table as
// it is. What the read answered with is recorded where the read is planned; see
// NormalizeOlapTableStreamScan and OlapTableStreamWrapper#answersWithTheCurrentTable.
if (rewriteContext.map(IvmRewriteContext::isReadFromAStreamOffset).orElse(false)) {
// Held rather than recorded: what these partitions hold is the image the read answered with,
// and the delta that follows is what brings the table up to date. See redeemHeldRecords.
LOG.info("Holding back the {} partitions of mv={}: their read answered with a base table as of "
+ "an older state than it is in now, taskId={}",
execPartitionNames.size(), mtmv.getName(), getTaskId());
holdCapturedEpochs(execPartitionSnapshots, batchCapturedEpochs);
} else {
partitionSnapshots.putAll(execPartitionSnapshots);
commitCapturedEpochs(batchCapturedEpochs);
}
}
if (capturePlanSignature) {
refreshedIvmPlanSignature = refreshedPlanSignature.getSha256();
}
LOG.info("MTMVTask refresh used snapshot: {}, mvDbName: {}, mvName: {}, taskId: {}", partitionSnapshots,
mtmv.getDatabase().getFullName(), mtmv.getName(), getTaskId());
}
/**
* Whether a base table the refresh reads has no stream that can be read, which no attempt other
* than COMPLETE can work around.
*
* <p>A base table that cannot be resolved is skipped rather than judged: it says nothing about the
* streams, and the refresh fails on it for its own reasons -- the attempt that runs reports that,
* this one only decides which attempt that should be.
*/
private boolean hasUnusableIvmStream() {
Database mvDb = (Database) mtmv.getDatabase();
if (mvDb == null) {
// Nothing to look the streams up in, so there is nothing to decide here.
return false;
}
Set<TableNameInfo> excluded = mtmv.getExcludedTriggerTables();
// The tables in the plan, not the relation's closure: a chained MV is created with a stream for
// every base table behind the MVs it reads, but no rewrite ever looks those up -- the incremental
// rewriter and the full refresh take the streams of the plan's scans -- so judging them would
// rebuild an MV whose refresh had nothing wrong with it.
for (BaseTableInfo baseTableInfo : relation.getBaseTablesOneLevelAndFromView()) {
OlapTable baseTable = resolveIvmBaseTable(baseTableInfo);
if (baseTable == null) {
continue;
}
if (MTMVPartitionUtil.isTableExcluded(excluded,
new TableNameInfo(baseTable.getFullQualifiers()))) {
continue;
}
if (usableIvmStream(mvDb, baseTable) == null) {
return true;
}
}
return false;
}
/**
* Whether a stream this partition refresh will read is missing or unusable.
*
* <p>The set is narrower than the plan, and which tables are in it depends on the partitions being
* refreshed: a PCT table that no refreshed partition's mapping names keeps its place in the plan as
* an ordinary scan, so its stream is never read, and a table outside the plan's own tables is not
* scanned at all. Judging the whole plan instead sends a partition refresh that would have worked to
* a full rebuild whenever such an unread stream is missing.
*/
private boolean hasUnusableIvmStreamForPartitions(MTMVRefreshContext context,
List<String> mvPartitionNames) {
Database mvDb = (Database) mtmv.getDatabase();
if (mvDb == null) {
// Nothing to look the streams up in, so there is nothing to decide here.
return false;
}
Set<TableNameInfo> excluded = mtmv.getExcludedTriggerTables();
Set<OlapTable> streamedTables = Sets.newLinkedHashSet();
Set<BaseTableInfo> pctTableInfos = Sets.newHashSet();
for (BaseColInfo pctInfo : mtmv.getMvPartitionInfo().getPctInfos()) {
pctTableInfos.add(pctInfo.getTableInfo());
}
// Every table of the plan that is not a PCT table is read through its stream: the partition
// refresh gives those tables a read mode (IvmRewriteContext.full).
for (BaseTableInfo baseTableInfo : relation.getBaseTablesOneLevelAndFromView()) {
if (pctTableInfos.contains(baseTableInfo)) {
continue;
}
OlapTable baseTable = resolveIvmBaseTable(baseTableInfo);
if (baseTable != null) {
streamedTables.add(baseTable);
}
}
// A PCT table is read through its stream only while the mapping of a refreshed partition names
// it, which is exactly the mapping the reset read mode is built from.
for (String mvPartitionName : mvPartitionNames) {
for (MTMVRelatedTableIf relatedTable : context.getByPartitionName(mvPartitionName).keySet()) {
if (relatedTable instanceof OlapTable) {
streamedTables.add((OlapTable) relatedTable);
}
}
}
for (OlapTable baseTable : streamedTables) {
if (MTMVPartitionUtil.isTableExcluded(excluded,
new TableNameInfo(baseTable.getFullQualifiers()))) {
continue;
}
if (usableIvmStream(mvDb, baseTable) == null) {
LOG.warn("IVM stream is unusable for the partitions this refresh plans, mv={}, baseTable={}",
mtmv.getName(), baseTable.getName());
return true;
}
}
return false;
}
/**
* A base table of the plan, or null when it cannot be resolved -- which says nothing about its stream,
* so a caller that is only judging streams skips it.
*/
private OlapTable resolveIvmBaseTable(BaseTableInfo baseTableInfo) {
try {
return (OlapTable) MTMVUtil.getTable(baseTableInfo);
} catch (Exception e) {
LOG.warn("Cannot resolve base table {} of mv={}", baseTableInfo, mtmv.getName(), e);
return null;
}
}
/** The MV's stream for the base table, or null when it is missing or cannot be used. */
private BaseTableStream usableIvmStream(Database mvDb, OlapTable baseTable) {
TableIf stream = mvDb.getTableNullable(
IvmUtil.streamName(mtmv.getId(), baseTable.getFullQualifiers()));
if (stream instanceof BaseTableStream
&& IvmUtil.isIvmStreamUsable((BaseTableStream) stream, baseTable)) {
return (BaseTableStream) stream;
}
return null;
}
private void reconcileIvmStreams(ConnectContext ctx) throws JobException {
try {
Database mvDb = (Database) mtmv.getDatabase();
Set<TableNameInfo> excluded = mtmv.getExcludedTriggerTables();
for (BaseTableInfo baseTableInfo : relation.getBaseTables()) {
OlapTable baseTable = (OlapTable) MTMVUtil.getTable(baseTableInfo);
if (MTMVPartitionUtil.isTableExcluded(excluded,
new TableNameInfo(baseTable.getFullQualifiers()))) {
continue;
}
if (usableIvmStream(mvDb, baseTable) != null) {
continue;
}
CreateMTMVCommand.createTableStream(ctx, mvDb, mtmv, baseTable);
}
} catch (UserException e) {
throw new JobException("Failed to reconcile IVM streams for mv=" + mtmv.getName(), e);
}
}
private Map<BaseTableInfo, Set<Long>> collectPctResetPartitionIds(MTMVRefreshContext context,
Set<String> execPartitionNames) throws AnalysisException {
Map<BaseTableInfo, Set<Long>> resetPartitionIds = Maps.newHashMap();
for (String mvPartitionName : execPartitionNames) {
for (Entry<MTMVRelatedTableIf, Set<String>> entry
: context.getByPartitionName(mvPartitionName).entrySet()) {
if (!(entry.getKey() instanceof OlapTable)) {
continue;
}
OlapTable pctTable = (OlapTable) entry.getKey();
Set<Long> partitionIds = resetPartitionIds.computeIfAbsent(
new BaseTableInfo(pctTable), key -> Sets.newHashSet());
for (String pctPartitionName : entry.getValue()) {
partitionIds.add(pctTable.getPartitionOrAnalysisException(pctPartitionName).getId());
}
}
}
return resetPartitionIds;
}
private IvmPlanSignature refreshPartitionsWithRetry(Set<String> execPartitionNames,
Map<TableIf, String> tableWithPartKey,
Optional<IvmRewriteContext> rewriteContext, RefreshMode refreshMode)
throws Exception {
return executeWithRetry(() -> refreshPartitions(execPartitionNames, tableWithPartKey,
rewriteContext, refreshMode),
"partition refresh, execPartitionNames=" + execPartitionNames);
}
private <T> T executeWithRetry(Callable<T> callable, String description) throws Exception {
int retryCount = 0;
int retryTime = Config.max_query_retry_time;
retryTime = retryTime <= 0 ? 1 : retryTime + 1;
Exception lastException = null;
while (retryCount < retryTime) {
if (TaskStatus.CANCELED.equals(getStatus())) {
throw new JobException("MTMV task is CANCELED");
}
try {
return callable.call();
} catch (Exception e) {
if (TaskStatus.CANCELED.equals(getStatus()) || !needRetry(e)) {
throw e;
}
lastException = e;
int randomMillis = 10 + (int) (Math.random() * 10);
if (retryCount > retryTime / 2) {
randomMillis = 20 + (int) (Math.random() * 10);
}
if (DebugPointUtil.isEnable("MTMVTask.retry.longtime")) {
randomMillis = 1000;
}
retryCount++;
LOG.warn("Retrying execution due to exception: {}. Attempt {}/{}, "
+ "taskId {} description {} lastQueryId {}, randomMillis {}",
e.getMessage(), retryCount, retryTime, getTaskId(),
description, lastQueryId, randomMillis);
if (retryCount >= retryTime) {
throw new Exception("Max retry attempts reached, original: " + lastException);
}
Thread.sleep(randomMillis);
}
}
throw new IllegalStateException("MTMV refresh retry loop exited unexpectedly");
}
private boolean needRetry(Exception exception) {
Throwable cause = exception;
while (cause != null) {
if (cause instanceof RpcException) {
return true;
}
cause = cause.getCause();
}
return Config.isCloudMode() && SystemInfoService.needRetryWithReplan(
Util.getRootCauseMessage(exception));
}
private IvmPlanSignature refreshPartitions(Set<String> refreshPartitionNames,
Map<TableIf, String> tableWithPartKey,
Optional<IvmRewriteContext> rewriteContext, RefreshMode refreshMode)
throws Exception {
// Create the MTMV context before parsing the MV definition SQL so SET_VAR hints
// resolve against the internal session (with MTMV disabled rules, etc.).
ConnectContext mtmvCtx = MTMVPlanUtil.createMTMVContext(mtmv, MTMVPlanUtil.DISABLE_RULES_WHEN_RUN_MTMV_TASK);
StatementContext statementContext = new StatementContext(
mtmvCtx, new OriginStatement(mtmv.getQuerySql(), 0));
// Install the StatementContext on the ConnectContext before parsing
// the MV definition SQL. UpdateMvByPartitionCommand.from() calls
// NereidsParser.parseSingle() which, for SQL containing SET_VAR hints,
// accesses ConnectContext.get().getStatementContext() inside
// LogicalPlanBuilder.withHints(). Without this assignment the
// StatementContext is null and a NullPointerException is thrown.
mtmvCtx.setStatementContext(statementContext);
statementContext.setConnectContext(mtmvCtx);
statementContext.setExcludedTriggerTables(mtmv.getExcludedTriggerTables());
statementContext.setIvmRewriteContext(rewriteContext);
for (Entry<MvccTableInfo, MvccSnapshot> entry : snapshots.entrySet()) {
statementContext.setSnapshot(entry.getKey(), entry.getValue());
}
// if SELF_MANAGE mv, only have default partition, will not have partitionItem, so we give empty set
UpdateMvByPartitionCommand command = UpdateMvByPartitionCommand
.from(mtmv, mtmv.getMvPartitionInfo().getPartitionType() != MTMVPartitionType.SELF_MANAGE
? refreshPartitionNames : Sets.newHashSet(), tableWithPartKey, statementContext);
setupComputeGroup(mtmvCtx);
AtomicReference<IvmPlanSignature> signatureRef = new AtomicReference<>();
try {
MTMVPlanUtil.executeCommand(mtmvCtx, command, statementContext,
getRefreshAuditStmt(refreshMode, refreshPartitionNames),
createRefreshConsumer(signatureRef));
} finally {
try {
recordQueryId(DebugUtil.printId(mtmvCtx.queryId()));
} finally {
closeExecutionContext(mtmvCtx);
}
}
if (getStatus() == TaskStatus.CANCELED) {
throw new JobException("task is CANCELED");
}
if (!rewriteContext.isPresent()) {
return null;
}
return Objects.requireNonNull(signatureRef.get(),
"IVM COMPLETE refresh did not produce a plan signature");
}
/**
* Builds the executor consumer for a COMPLETE refresh: registers the executing
* statement for cancellation, and when the command finishes (consumer invoked with
* {@code null}), captures the IVM plan signature into {@code signatureRef} from the
* still-registered executor before it is cleared.
*/
private Consumer<StmtExecutor> createRefreshConsumer(AtomicReference<IvmPlanSignature> signatureRef) {
return executor -> {
// executor == null ��� this.executor ���������:executeCommand ��� finally ������,
// ���������������������,��� registerExecutor(null) ������������������������ IVM ���������
if (executor == null && this.executor != null && this.executor.planner() != null) {
((NereidsPlanner) this.executor.planner())
.getCascadesContext()
.getIvmRewriteResult()
.ifPresent(r -> signatureRef.set(r.getPlanSignature()));
}
registerExecutor(executor);
};
}
private static void closeExecutionContext(ConnectContext executionContext) {
try {
if (executionContext.queryId() != null) {
QeProcessorImpl.INSTANCE.unregisterQuery(executionContext.queryId());
}
} finally {
StatementContext statementContext = executionContext.getStatementContext();
if (statementContext != null) {
statementContext.close();
}
}
}
private void setupComputeGroup(ConnectContext ctx) {
if (!Config.isCloudMode()) {
computeGroup = FeConstants.null_string;
return;
}
String taskComputeGroup = taskContext.getComputeGroup();
if (!Strings.isNullOrEmpty(taskComputeGroup)) {
ctx.setCloudCluster(taskComputeGroup);
}
try {
computeGroup = ctx.getCloudCluster(false);
} catch (ComputeGroupException e) {
computeGroup = FeConstants.null_string;
LOG.warn("failed to resolve compute group for mtmv task, taskId: {}", getTaskId(), e);
}
}
private void recordQueryId(String queryId) {
if (!Strings.isNullOrEmpty(queryId)) {
lastQueryId = queryId;
}
}
/**
* Registers the currently executing statement (or clears it with {@code null}).
* Called by the executeCommand executor consumer: non-null right before the command
* runs, null after it finishes. The cancel path reads this to interrupt execution.
* Rejects registration with an exception when the task has already been cancelled,
* so the command is not started at all. Uses an unchecked exception so the method
* reference can be used as a {@code Consumer<StmtExecutor>}.
*/
public void registerExecutor(StmtExecutor executor) {
if (executor != null && getStatus() == TaskStatus.CANCELED) {
this.executor = null;
throw new IllegalStateException("task is CANCELED");
}
this.executor = executor;
}
private String getRefreshAuditStmt(RefreshMode refreshMode, Set<String> refreshPartitionNames) {
String mvName = mtmv.getName();
DatabaseIf database = mtmv.getDatabase();
if (database != null) {
mvName = database.getFullName() + "." + mvName;
CatalogIf catalog = database.getCatalog();
if (catalog != null) {
mvName = catalog.getName() + "." + mvName;
}
}
return String.format(
"Asynchronous materialized view refresh task, mvName: %s,"
+ "taskId: %s, refreshMode: %s, partitions: %s",
mvName, super.getTaskId(), refreshMode, refreshPartitionNames);
}
@Override
public synchronized boolean onFail() throws JobException {
LOG.info("mtmv task onFail, taskId: {}", super.getTaskId());
boolean res = super.onFail();
if (!res) {
return false;
}
after();
if (MetricRepo.isInit) {
MetricRepo.COUNTER_ASYNC_MATERIALIZED_VIEW_TASK_FAILED_NUM.increase(1L);
}
return true;
}
@Override
public synchronized boolean onSuccess() throws JobException {
if (LOG.isDebugEnabled()) {
LOG.debug("mtmv task onSuccess, taskId: {}", super.getTaskId());
}
boolean res = super.onSuccess();
if (!res) {
return false;
}
after();
if (MetricRepo.isInit) {
MetricRepo.HISTO_ASYNC_MATERIALIZED_VIEW_TASK_DURATION.update(
super.getFinishTimeMs() - super.getStartTimeMs());
MetricRepo.COUNTER_ASYNC_MATERIALIZED_VIEW_TASK_SUCCESS_NUM.increase(1L);
}
return true;
}
/**
* The reason for overriding the parent class is to add synchronized protection
*/
@Override
public synchronized boolean cancel(boolean needWaitCancelComplete) throws JobException {
if (LOG.isDebugEnabled()) {
LOG.debug("mtmv task cancel, taskId: {}", super.getTaskId());
}
return super.cancel(needWaitCancelComplete);
}
@Override
protected void executeCancelLogic(boolean needWaitCancelComplete) {
try {
// Mtmv is initialized in the before method.
// If the task has not yet run, the before method will not be used, so mtmv will be empty,
// which prevents the canceled task from being added to the history list
if (mtmv == null) {
mtmv = MTMVUtil.getMTMV(dbId, mtmvId);
}
} catch (UserException e) {
LOG.warn("executeCancelLogic failed:", e);
return;
}
if (executor != null) {
executor.cancel(new Status(TStatusCode.CANCELLED, "mtmv task cancelled"), needWaitCancelComplete);
}
after();
}
@Override
public void before() throws JobException {
if (LOG.isDebugEnabled()) {
LOG.debug("mtmv task before, taskId: {}", super.getTaskId());
}
super.before();
try {
mtmv = MTMVUtil.getMTMV(dbId, mtmvId);
} catch (UserException e) {
LOG.warn("before task failed:", e);
throw new JobException(e);
}
}
/**
* Do something before refreshing, such as clearing the cache of the external table
*
* @throws AnalysisException
* @throws DdlException
*/
private void beforeMTMVRefresh() throws AnalysisException, DdlException {
for (BaseTableInfo tableInfo : relation.getBaseTablesOneLevelAndFromView()) {
TableIf tableIf = MTMVUtil.getTable(tableInfo);
if (tableIf instanceof MTMVBaseTableIf) {
MTMVBaseTableIf baseTableIf = (MTMVBaseTableIf) tableIf;
baseTableIf.beforeMTMVRefresh(mtmv);
}
if (tableIf instanceof MvccTable) {
MvccTable mvccTable = (MvccTable) tableIf;
MvccSnapshot mvccSnapshot = mvccTable.loadSnapshot(Optional.empty(), Optional.empty());
snapshots.put(new MvccTableInfo(mvccTable), mvccSnapshot);
}
}
}
@Override
public void runTask() throws JobException {
LOG.info("mtmv task runTask, taskId: {}", super.getTaskId());
MTMVJob job = (MTMVJob) getJobOrJobException();
try {
if (LOG.isDebugEnabled()) {
LOG.debug("mtmv task get writeLock start, taskId: {}", super.getTaskId());
}
job.writeLock();
if (LOG.isDebugEnabled()) {
LOG.debug("mtmv task get writeLock end, taskId: {}", super.getTaskId());
}
super.runTask();
} finally {
job.writeUnlock();
if (LOG.isDebugEnabled()) {
LOG.debug("mtmv task release writeLock, taskId: {}", super.getTaskId());
}
}
}
@Override
public TRow getTvfInfo(String jobName) {
TRow trow = new TRow();
trow.addToColumnValue(new TCell().setStringVal(String.valueOf(super.getTaskId())));
trow.addToColumnValue(new TCell().setStringVal(String.valueOf(super.getJobId())));
trow.addToColumnValue(new TCell().setStringVal(jobName));
String dbName = "";
String mvName = "";
try {
MTMV mtmv = MTMVUtil.getMTMV(dbId, mtmvId);
dbName = mtmv.getQualifiedDbName();
mvName = mtmv.getName();
} catch (UserException e) {
LOG.warn("can not find mv", e);
}
trow.addToColumnValue(new TCell().setStringVal(String.valueOf(mtmvId)));
trow.addToColumnValue(new TCell().setStringVal(mvName));
trow.addToColumnValue(new TCell().setStringVal(String.valueOf(dbId)));
trow.addToColumnValue(new TCell().setStringVal(dbName));
trow.addToColumnValue(new TCell()
.setStringVal(super.getStatus() == null ? FeConstants.null_string : super.getStatus().toString()));
trow.addToColumnValue(new TCell().setStringVal(super.getErrMsg()));
trow.addToColumnValue(new TCell().setStringVal(TimeUtils.longToTimeString(super.getCreateTimeMs())));
trow.addToColumnValue(new TCell().setStringVal(TimeUtils.longToTimeString(super.getStartTimeMs())));
trow.addToColumnValue(new TCell().setStringVal(TimeUtils.longToTimeString(super.getFinishTimeMs())));
trow.addToColumnValue(new TCell().setStringVal(
(super.getFinishTimeMs() == null || super.getFinishTimeMs() == 0) ? FeConstants.null_string
: String.valueOf(super.getFinishTimeMs() - super.getStartTimeMs())));
trow.addToColumnValue(new TCell()
.setStringVal(taskContext == null ? FeConstants.null_string : GSON.toJson(taskContext)));
trow.addToColumnValue(
new TCell().setStringVal(refreshMode == null ? FeConstants.null_string : refreshMode.toString()));
trow.addToColumnValue(
new TCell().setStringVal(
needRefreshPartitions == null ? FeConstants.null_string : GSON.toJson(
needRefreshPartitions)));
trow.addToColumnValue(
new TCell().setStringVal(
completedPartitions == null ? FeConstants.null_string : GSON.toJson(
completedPartitions)));
trow.addToColumnValue(
new TCell().setStringVal(getProgress()));
trow.addToColumnValue(
new TCell().setStringVal(lastQueryId));
trow.addToColumnValue(new TCell().setStringVal(
computeGroup == null || computeGroup.isEmpty() ? FeConstants.null_string : computeGroup));
trow.addToColumnValue(new TCell().setStringVal(
ivmFallbackReason == null ? FeConstants.null_string : ivmFallbackReason));
trow.addToColumnValue(new TCell().setStringVal(String.valueOf(ivmRebuiltPartitions)));
return trow;
}
private String getProgress() {
if (CollectionUtils.isEmpty(needRefreshPartitions)) {
return FeConstants.null_string;
}
int completedSize = CollectionUtils.isEmpty(completedPartitions) ? 0 : completedPartitions.size();
BigDecimal result = new BigDecimal(completedSize * 100)
.divide(new BigDecimal(needRefreshPartitions.size()), 2, RoundingMode.HALF_UP);
StringBuilder builder = new StringBuilder(result.toString());
builder.append("% (");
builder.append(completedSize);
builder.append("/");
builder.append(needRefreshPartitions.size());
builder.append(")");
return builder.toString();
}
private void after() {
if (mtmv != null) {
// A copy, because this is not always the thread that fills the map: a STOP publishes from the
// cancel thread and `cancel(false)` does not wait for the execution to stop, so the batches the
// worker commits afterwards would otherwise land in a result that was already built -- and
// journaled, since the record is written out asynchronously -- as theirs. See
// getIvmCapturedEpochs for the same rule on the epochs.
Map<String, MTMVRefreshPartitionSnapshot> publishedSnapshots = partitionSnapshots == null
? null : Maps.newHashMap(partitionSnapshots);
Env.getCurrentEnv()
.addMTMVTaskResult(new TableNameInfo(mtmv.getQualifiedDbName(), mtmv.getName()), this, relation,
publishedSnapshots);
}
}
@Override
protected void closeOrReleaseResources() {
if (null != mtmv) {
mtmv = null;
}
if (null != executor) {
executor = null;
}
if (null != relation) {
relation = null;
}
if (null != partitionSnapshots) {
partitionSnapshots = null;
}
if (null != snapshots) {
snapshots = null;
}
}
private Map<TableIf, String> getIncrementalTableMap() throws AnalysisException {
Map<TableIf, String> tableWithPartKey = Maps.newHashMap();
if (mtmv.getMvPartitionInfo().getPartitionType() != MTMVPartitionType.SELF_MANAGE) {
List<BaseColInfo> pctInfos = mtmv.getMvPartitionInfo().getPctInfos();
for (BaseColInfo pctInfo : pctInfos) {
tableWithPartKey
.put(MTMVUtil.getTable(pctInfo.getTableInfo()), pctInfo.getColName());
}
}
return tableWithPartKey;
}
private MTMVTaskRefreshMode generateRefreshMode(List<String> needRefreshPartitionIds) {
if (CollectionUtils.isEmpty(needRefreshPartitionIds)) {
return MTMVTaskRefreshMode.NOT_REFRESH;
} else if (Sets.newHashSet(needRefreshPartitionIds).equals(mtmv.getPartitionNames())) {
return MTMVTaskRefreshMode.COMPLETE;
} else {
return MTMVTaskRefreshMode.PARTIAL;
}
}
public MTMVTaskContext getTaskContext() {
return taskContext;
}
public long getMtmvSchemaChangeVersion() {
return mtmvSchemaChangeVersion;
}
/**
* The epochs this task's committed batches were read at -- the requirement each refreshed partition was
* read under, see {@link #captureLatestEpochs} -- as a detached copy.
*
* <p>Detached rather than live: the batches merge into the map on the thread running the task, while a
* STOP publishes what it holds from the callback thread -- `cancel(false)` does not wait for the
* execution to stop -- so a caller reading the field itself would iterate a map still being written, and
* would journal a result that goes on changing after it was read.
*/
public Map<String, Long> getIvmCapturedEpochs() {
// A task read back from the journal has none: the field is transient, so gson leaves it null and the
// constructor that would have initialized it never runs. What a replay applies is the states the
// record carries, not this.
return ivmCapturedEpochs == null ? Collections.emptyMap() : Maps.newHashMap(ivmCapturedEpochs);
}
public String getRefreshedIvmPlanSignature() {
return refreshedIvmPlanSignature;
}
@Override
public String toString() {
return "MTMVTask{"
+ "dbId=" + dbId
+ ", mtmvId=" + mtmvId
+ ", taskContext=" + taskContext
+ ", needRefreshPartitions=" + needRefreshPartitions
+ ", completedPartitions=" + completedPartitions
+ ", refreshMode=" + refreshMode
+ "} " + super.toString();
}
}