InsertOverwriteTableCommand.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.nereids.trees.plans.commands.insert;
import org.apache.doris.analysis.StmtType;
import org.apache.doris.catalog.Env;
import org.apache.doris.catalog.MTMV;
import org.apache.doris.catalog.OlapTable;
import org.apache.doris.catalog.TableIf;
import org.apache.doris.common.ErrorCode;
import org.apache.doris.common.ErrorReport;
import org.apache.doris.common.UserException;
import org.apache.doris.common.util.DebugPointUtil;
import org.apache.doris.common.util.InternalDatabaseUtil;
import org.apache.doris.connector.spi.handle.WriteOperation;
import org.apache.doris.datasource.doris.RemoteDorisExternalTable;
import org.apache.doris.datasource.doris.RemoteOlapTable;
import org.apache.doris.datasource.plugin.PluginDrivenExternalTable;
import org.apache.doris.insertoverwrite.AbstractInsertOverwriteManager;
import org.apache.doris.insertoverwrite.InsertOverwriteUtil;
import org.apache.doris.insertoverwrite.RemoteInsertOverwriteManager;
import org.apache.doris.mtmv.MTMVUtil;
import org.apache.doris.mysql.privilege.PrivPredicate;
import org.apache.doris.nereids.CascadesContext;
import org.apache.doris.nereids.NereidsPlanner;
import org.apache.doris.nereids.StatementContext;
import org.apache.doris.nereids.analyzer.UnboundConnectorTableSink;
import org.apache.doris.nereids.analyzer.UnboundTableSink;
import org.apache.doris.nereids.analyzer.UnboundTableSinkCreator;
import org.apache.doris.nereids.exceptions.AnalysisException;
import org.apache.doris.nereids.glue.LogicalPlanAdapter;
import org.apache.doris.nereids.lineage.LineageInfoExtractor;
import org.apache.doris.nereids.lineage.LineageUtils;
import org.apache.doris.nereids.properties.PhysicalProperties;
import org.apache.doris.nereids.trees.TreeNode;
import org.apache.doris.nereids.trees.expressions.Expression;
import org.apache.doris.nereids.trees.expressions.literal.Literal;
import org.apache.doris.nereids.trees.plans.Explainable;
import org.apache.doris.nereids.trees.plans.Plan;
import org.apache.doris.nereids.trees.plans.PlanType;
import org.apache.doris.nereids.trees.plans.algebra.TVFRelation;
import org.apache.doris.nereids.trees.plans.commands.Command;
import org.apache.doris.nereids.trees.plans.commands.ForwardWithSync;
import org.apache.doris.nereids.trees.plans.commands.NeedAuditEncryption;
import org.apache.doris.nereids.trees.plans.logical.LogicalPlan;
import org.apache.doris.nereids.trees.plans.logical.UnboundLogicalSink;
import org.apache.doris.nereids.trees.plans.physical.PhysicalOlapTableSink;
import org.apache.doris.nereids.trees.plans.physical.PhysicalTableSink;
import org.apache.doris.nereids.trees.plans.visitor.PlanVisitor;
import org.apache.doris.planner.ScanNode;
import org.apache.doris.qe.ConnectContext;
import org.apache.doris.qe.QueryState.MysqlStateType;
import org.apache.doris.qe.StmtExecutor;
import org.apache.doris.thrift.TPartialUpdateNewRowPolicy;
import com.google.common.base.Preconditions;
import com.google.common.base.Strings;
import com.google.common.collect.Lists;
import com.google.common.collect.Maps;
import org.apache.commons.collections4.CollectionUtils;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.awaitility.Awaitility;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
/**
* insert into select command implementation
* insert into select command support the grammer: explain? insert into table columns? partitions? hints? query
* InsertIntoTableCommand is a command to represent insert the answer of a query into a table.
* class structure's:
* InsertIntoTableCommand(Query())
* ExplainCommand(Query())
*/
public class InsertOverwriteTableCommand extends Command
implements NeedAuditEncryption, ForwardWithSync, Explainable, CancelableCommand {
/**
* Fails an overwrite in the one window a refresh cannot recover from by itself: after the rows have been
* committed into the temporary partitions and before the swap publishes them. Everything the write read
* is committed by then -- the base table streams it consumed, among them -- and the partitions it was
* going to replace still hold what they had, so a refresh that dies here leaves rows missing and nothing
* durable saying so unless it raised a rebuild requirement before it read. See
* test_ivm_overwrite_failure_between_the_halves, which pins the recovery.
*
* <p>Scoped by the MV name the point carries as its {@code mv_name} parameter: the read below answers
* with the default when the point is not enabled or carries no such parameter, and no MV is named by an
* empty string, so enabling it cannot disturb an overwrite that is not the one under test.
*/
public static final String DEBUG_POINT_FAIL_BETWEEN_THE_HALVES_OF_AN_OVERWRITE =
"InsertOverwriteTableCommand.failBetweenTheTwoHalvesOfAnOverwrite";
/**
* Cancels the overwrite before it has committed anything, so that the half of a cancellation's meaning
* that takes the statement back can be pinned by a test: nothing durable happened, so the statement has
* to fail rather than report the success of an overwrite that did not run. See test_insert_overwrite_cancel.
*
* <p>Its {@code table_name} parameter names the one table the point may disturb: the read below answers
* with the default when the point is not enabled or carries no such parameter, and no table is named by an
* empty string, so an enabled point cannot disturb an overwrite that is not the one under test.
*/
public static final String DEBUG_POINT_CANCEL_BEFORE_THE_INSERT_OF_AN_OVERWRITE =
"InsertOverwriteTableCommand.cancelBeforeTheInsertOfAnOverwrite";
/**
* Cancels the overwrite in the window between its two halves -- after the insert, before the swap -- so
* that the other half of a cancellation's meaning can be pinned: where the rows are durable, the swap runs
* and the statement reports the success it is; where the insert committed nothing, the cancellation still
* has everything to take back. See test_insert_overwrite_cancel.
*
* <p>Separate from the point above rather than one point with a stage parameter: a point is consumed by
* the first lookup that reads it (see {@code DebugPointUtil#getDebugPoint}), so a shared name would let
* the check at one site spend the other site's allowance and make an armed point silently not fire.
*/
public static final String DEBUG_POINT_CANCEL_BETWEEN_THE_HALVES_OF_AN_OVERWRITE =
"InsertOverwriteTableCommand.cancelBetweenTheTwoHalvesOfAnOverwrite";
/**
* Cancels the overwrite while the swap holds the target table's write lock, which is the window a
* cancellation can reach only after the check that reads the flag before the swap was issued: the swap
* waits for that lock, and the wait can be as long as whoever holds it. A cancellation with nothing
* committed is still honoured there, because there is nothing durable to publish and refusing costs the
* statement and nothing else. See test_insert_overwrite_cancel.
*/
public static final String DEBUG_POINT_CANCEL_WHILE_THE_SWAP_WAITS_FOR_THE_TABLE_LOCK =
"InsertOverwriteTableCommand.cancelWhileTheSwapWaitsForTheTableLock";
/**
* The swap that publishes an overwrite: replacing the temp partitions for an explicit-partition
* overwrite, or making a task group's replacements visible for an auto-detect one. Both run through
* {@link #publishTheOverwrite}, which holds the target table's write lock and takes the last look at the
* cancellation flag before letting one run.
*/
@FunctionalInterface
interface OverwritePublication {
void publish() throws UserException;
}
private static final Logger LOG = LogManager.getLogger(InsertOverwriteTableCommand.class);
private LogicalPlan originLogicalQuery;
private Optional<LogicalPlan> logicalQuery;
private Optional<String> labelName;
private final Optional<LogicalPlan> cte;
private AtomicBoolean isCancelled = new AtomicBoolean(false);
private AtomicBoolean isRunning = new AtomicBoolean(false);
private Optional<String> branchName;
private Optional<Plan> lineagePlan = Optional.empty();
/**
* constructor
*/
public InsertOverwriteTableCommand(LogicalPlan logicalQuery, Optional<String> labelName,
Optional<LogicalPlan> cte, Optional<String> branchName) {
super(PlanType.INSERT_INTO_TABLE_COMMAND);
this.originLogicalQuery = Objects.requireNonNull(logicalQuery, "logicalQuery should not be null");
this.logicalQuery = Optional.empty();
this.labelName = Objects.requireNonNull(labelName, "labelName should not be null");
this.cte = cte;
this.branchName = branchName;
}
public void setLabelName(Optional<String> labelName) {
this.labelName = labelName;
}
public boolean isAutoDetectOverwrite(LogicalPlan logicalQuery) {
return (logicalQuery instanceof UnboundTableSink)
&& ((UnboundTableSink<?>) logicalQuery).isAutoDetectPartition();
}
public LogicalPlan getLogicalQuery() {
return logicalQuery.orElse(originLogicalQuery);
}
@Override
public void run(ConnectContext ctx, StmtExecutor executor) throws Exception {
TableIf targetTableIf = InsertUtils.getTargetTable(originLogicalQuery, ctx);
// check allow insert overwrite
if (!allowInsertOverwrite(targetTableIf)) {
String errMsg = "insert into overwrite only support OLAP/Remote OLAP table and external"
+ " tables (HMS/Iceberg, or a plugin-driven connector that supports overwrite)."
+ " But current table type is " + targetTableIf.getType();
LOG.error(errMsg);
throw new AnalysisException(errMsg);
}
//check allow modify MTMVData
if (targetTableIf instanceof MTMV && !MTMVUtil.allowModifyMTMVData(ctx)) {
throw new AnalysisException("Not allowed to perform current operation on async materialized view");
}
// Check the branch capability before resolving the branch-specific writer schema. Otherwise,
// an unsupported connector can fail while resolving a branch instead of reporting the
// INSERT OVERWRITE capability error.
if (branchName.isPresent() && !pluginConnectorSupportsWriteBranch(targetTableIf)) {
throw new AnalysisException(
"Only support insert overwrite into iceberg table's branch");
}
ctx.getStatementContext().setIsInsert(true);
Optional<CascadesContext> analyzeContext = Optional.of(
CascadesContext.initContext(ctx.getStatementContext(), originLogicalQuery, PhysicalProperties.ANY)
);
InsertUtils.pinConnectorWriteSchema(
ctx.getStatementContext(), targetTableIf, originLogicalQuery, branchName);
this.logicalQuery = Optional.of((LogicalPlan) InsertUtils.normalizePlan(
originLogicalQuery, (targetTableIf instanceof RemoteDorisExternalTable)
? ((RemoteDorisExternalTable) targetTableIf).getOlapTable() : targetTableIf,
analyzeContext, Optional.empty()));
if (cte.isPresent()) {
LogicalPlan logicalQuery = this.logicalQuery.get();
this.logicalQuery = Optional.of(
(LogicalPlan) logicalQuery.withChildren(
cte.get().withChildren(logicalQuery.child(0))
)
);
}
LogicalPlan logicalQuery = this.logicalQuery.get();
LogicalPlanAdapter logicalPlanAdapter = new LogicalPlanAdapter(logicalQuery, ctx.getStatementContext());
NereidsPlanner planner = new NereidsPlanner(ctx.getStatementContext());
LineageInfoExtractor.registerAnalyzePlanHook(ctx.getStatementContext(), planner);
planner.plan(logicalPlanAdapter, ctx.getSessionVariable().toThrift());
// This plan only locates the sink and the partitions; the insert below plans again and runs
// that plan. No coordinator ever takes this one, so what its scan nodes opened for the
// backend while planning (a remote Doris scan's Flight SQL session on the other frontend, a
// batch split source) is released here, before the real insert opens its own.
for (ScanNode scanNode : planner.getScanNodes()) {
scanNode.stop();
}
Plan analyzedPlan = planner.getAnalyzedPlan();
lineagePlan = Optional.ofNullable(analyzedPlan);
executor.checkBlockRules();
Optional<TreeNode<?>> plan = (planner.getPhysicalPlan()
.<TreeNode<?>>collect(node -> node instanceof PhysicalTableSink)).stream().findAny();
Preconditions.checkArgument(plan.isPresent(), "insert into command must contain OlapTableSinkNode");
PhysicalTableSink<?> physicalTableSink = ((PhysicalTableSink<?>) plan.get());
TableIf targetTable = physicalTableSink.getTargetTable();
List<String> partitionNames;
boolean wholeTable = false;
if (physicalTableSink instanceof PhysicalOlapTableSink) {
if (targetTable instanceof OlapTable) {
InternalDatabaseUtil
.checkDatabase(((OlapTable) targetTable).getQualifiedDbName(), ConnectContext.get());
// check auth
if (!Env.getCurrentEnv().getAccessManager()
.checkTblPriv(ConnectContext.get(), targetTable.getDatabase().getCatalog().getName(),
((OlapTable) targetTable).getQualifiedDbName(),
targetTable.getName(), PrivPredicate.LOAD)) {
ErrorReport.reportAnalysisException(ErrorCode.ERR_TABLEACCESS_DENIED_ERROR, "LOAD",
ConnectContext.get().getQualifiedUser(), ConnectContext.get().getRemoteIP(),
((OlapTable) targetTable).getQualifiedDbName() + ": " + targetTable.getName());
}
}
partitionNames = ((UnboundTableSink<?>) logicalQuery).getPartitions();
// If not specific partition to overwrite, means it's a command to overwrite the table.
// not we execute as overwrite every partitions.
if (CollectionUtils.isEmpty(partitionNames)) {
wholeTable = true;
try { // avoid concurrent modification exception when get partition names
targetTable.readLock();
partitionNames = Lists.newArrayList(targetTable.getPartitionNames());
} finally {
targetTable.readUnlock();
}
}
} else {
// Do not create temp partition on FE
partitionNames = new ArrayList<>();
}
AbstractInsertOverwriteManager insertOverwriteManager = (targetTable instanceof RemoteOlapTable)
? new RemoteInsertOverwriteManager(((RemoteOlapTable) targetTable).getCatalog())
: Env.getCurrentEnv().getInsertOverwriteManager();
insertOverwriteManager.recordRunningTableOrException(targetTable.getDatabase(), targetTable);
isRunning.set(true);
long taskId = 0;
try {
// OLAP overwrite runs its internal partition replacement with the auth check skipped.
// Set the flag here, inside the try, so the finally below always pairs the reset even if
// an earlier step (e.g. the @branch guard) throws before we get here.
if (physicalTableSink instanceof PhysicalOlapTableSink && targetTable instanceof OlapTable) {
ctx.setSkipAuth(true);
}
if (isAutoDetectOverwrite(getLogicalQuery())) {
// taskId here is a group id. it contains all replace tasks made and registered in rpc process.
final long groupId = insertOverwriteManager.registerTaskGroup(targetTable);
taskId = groupId;
// When inserting, BE will call to replace partition by FrontendService. FE will register new temp
// partitions and return. for transactional, the replacement will really occur when insert successed,
// i.e. `insertInto` finished. then we call taskGroupSuccess to make replacement.
InsertCommandContext insertCtx = insertIntoAutoDetect(ctx, executor, groupId);
if (isCancelled.get() && !insertCtx.hasCommitted()) {
// The load committed nothing -- an empty plan takes the path that begins no transaction --
// so the cancellation still has everything to take back: the catch drops the group's temp
// partitions (an empty plan registers none), and the statement fails rather than reporting
// a replacement the client cancelled. A cancellation landing after this check, while the
// swap waits for the table lock, is taken up again by publishTheOverwrite below.
throw cancelledBeforeTheRowsWereCommitted("after a load that committed nothing", ctx);
}
// The replacement stays under the lock publishTheOverwrite takes, and the group's bookkeeping
// follows it: its edit-log writes wait for their journals, and the target's readers and
// writers must not be blocked through them.
publishTheOverwrite(targetTable, insertCtx, ctx,
() -> insertOverwriteManager.replacePartitionsOfTaskGroup(groupId, (OlapTable) targetTable,
isForceDropPartition()));
insertOverwriteManager.finishTaskGroup(groupId);
} else {
// it's overwrite table(as all partitions) or specific partition(s)
List<String> tempPartitionNames = InsertOverwriteUtil.generateTempPartitionNames(partitionNames);
cancelTheOverwriteAt(DEBUG_POINT_CANCEL_BEFORE_THE_INSERT_OF_AN_OVERWRITE, targetTable);
if (isCancelled.get()) {
// Nothing durable happened: no task is registered, no temp partition exists, no row was
// written and nothing was committed. The statement is a plain failure, like the one the
// inner insert reports when it is cancelled, rather than the success of an overwrite that
// did not run.
throw cancelledBeforeTheRowsWereCommitted("before registerTask", ctx);
}
taskId = insertOverwriteManager.registerTask(targetTable, tempPartitionNames);
if (isCancelled.get()) {
// The catch below takes the registration back; no temp partition exists yet, so there is
// nothing else to drop.
throw cancelledBeforeTheRowsWereCommitted("before addTempPartitions", ctx);
}
InsertOverwriteUtil.addTempPartitions(targetTable, partitionNames, tempPartitionNames);
if (isCancelled.get()) {
// The catch below drops the temp partitions this cancelled statement created.
throw cancelledBeforeTheRowsWereCommitted("before insertInto", ctx);
}
// todo: need to refresh remote target table after add temp partitions
InsertCommandContext insertCtx = insertIntoPartitions(ctx, executor, tempPartitionNames, wholeTable);
cancelTheOverwriteAt(DEBUG_POINT_CANCEL_BETWEEN_THE_HALVES_OF_AN_OVERWRITE, targetTable);
if (isCancelled.get()) {
if (!insertCtx.hasCommitted()) {
// The insert committed nothing: its plan folded to an empty relation, so it took the
// path that begins no transaction, and this window holds no durable work at all.
// Completing the swap would publish an empty table for a statement the client
// cancelled; the catch drops the empty temp partitions instead and the statement
// fails, which is the same boundary the cancellations above sit on.
throw cancelledBeforeTheRowsWereCommitted("after an insert that committed nothing", ctx);
}
// Too late to cancel: insertIntoPartitions returns only once its transaction has committed
// the rows into the temp partitions -- visible, or still waiting for a publication that
// timed out -- and everything the read consumed, the base table stream offsets among it,
// was committed with that same transaction. Dropping the temp partitions here is exactly
// what would lose those rows against an advanced offset, while the swap below is what
// publishes them. The overwrite completes, and it is the outcome the statement reports.
LOG.info("insert overwrite is cancelled after its rows were committed, completing it,"
+ " queryId: {}", ctx.getQueryIdentifier());
}
failBetweenTheTwoHalvesOfAnOverwrite(targetTable);
// The publication below is a lambda, so the partitions it replaces need a name that is final:
// partitionNames is assigned on more than one path above.
final List<String> replacedPartitionNames = partitionNames;
publishTheOverwrite(targetTable, insertCtx, ctx,
() -> InsertOverwriteUtil.replacePartition(targetTable, replacedPartitionNames,
tempPartitionNames, isForceDropPartition()));
if (isCancelled.get()) {
LOG.info("insert overwrite is cancelled before taskSuccess, do nothing, queryId: {}",
ctx.getQueryIdentifier());
}
insertOverwriteManager.taskSuccess(taskId);
}
} catch (Exception e) {
LOG.warn("insert into overwrite failed with task(or group) id {}", taskId, e);
// A cancel that landed before registerTask leaves nothing registered to fail, and no id was taken.
if (isAutoDetectOverwrite(getLogicalQuery()) && taskId != 0) {
insertOverwriteManager.taskGroupFail(taskId);
} else if (taskId != 0) {
insertOverwriteManager.taskFail(taskId);
}
throw e;
} finally {
ConnectContext.get().setSkipAuth(false);
insertOverwriteManager.dropRunningRecord(targetTable.getDatabase(), targetTable);
isRunning.set(false);
}
LineageUtils.submitLineageEventIfNeeded(executor, lineagePlan, getLogicalQuery(), getClass());
}
/**
* cancel insert overwrite
*/
public void cancel() {
this.isCancelled.set(true);
}
/**
* wait insert overwrite not running
*/
public void waitNotRunning() {
long waitMaxTimeSecond = 10L;
try {
Awaitility.await().atMost(waitMaxTimeSecond, TimeUnit.SECONDS).untilFalse(isRunning);
} catch (Exception e) {
LOG.warn("waiting time exceeds {} second, stop wait, labelName: {}", waitMaxTimeSecond,
labelName.isPresent() ? labelName.get() : "", e);
}
}
private boolean allowInsertOverwrite(TableIf targetTable) {
if (targetTable instanceof OlapTable || targetTable instanceof RemoteDorisExternalTable) {
return true;
} else {
return targetTable instanceof PluginDrivenExternalTable
&& pluginConnectorSupportsInsertOverwrite((PluginDrivenExternalTable) targetTable);
}
}
/**
* A plugin-driven (SPI connector) table supports INSERT OVERWRITE only if its connector
* declares the capability. Connectors that support plain INSERT but not overwrite (e.g. jdbc)
* must be rejected here so the command fails loud, rather than reaching the sink and silently
* degrading OVERWRITE to a plain append. Mirrors the connector-access pattern in
* {@code PhysicalPlanTranslator}.
*/
private static boolean pluginConnectorSupportsInsertOverwrite(PluginDrivenExternalTable table) {
// Per-handle write-op probe (a heterogeneous gateway answers per-table; OVERWRITE happens to be admitted
// by both hive and iceberg, but the probe is resolved uniformly with the other write-op admission gates).
return table.connectorSupportedWriteOperations().contains(WriteOperation.OVERWRITE);
}
/**
* A plugin-driven (SPI connector) table accepts an {@code INSERT OVERWRITE t@branch(name)} only if
* its connector declares {@code supportsWriteBranch()}. Connectors with no branch concept must be
* rejected here (fail loud) instead of reaching the generic sink, which would silently drop the
* branch and overwrite the table's default ref. Mirrors {@code pluginConnectorSupportsInsertOverwrite}.
*/
private static boolean pluginConnectorSupportsWriteBranch(TableIf targetTable) {
if (!(targetTable instanceof PluginDrivenExternalTable)) {
return false;
}
// Per-handle: a heterogeneous gateway supports write-to-branch for its iceberg tables but not its hive.
return ((PluginDrivenExternalTable) targetTable).connectorSupportsWriteBranch();
}
/**
* Throws when the debug point names the MV this overwrite targets; see the constant above.
*/
private static void failBetweenTheTwoHalvesOfAnOverwrite(TableIf targetTable) throws UserException {
if (!(targetTable instanceof MTMV)
|| !targetTable.getName().equals(DebugPointUtil.getDebugParamOrDefault(
DEBUG_POINT_FAIL_BETWEEN_THE_HALVES_OF_AN_OVERWRITE, "mv_name", ""))) {
return;
}
throw new UserException("debug point: " + DEBUG_POINT_FAIL_BETWEEN_THE_HALVES_OF_AN_OVERWRITE);
}
/**
* Cancels this overwrite when the debug point names the table it targets; see the constants above.
* Nothing here decides what a cancelled overwrite means -- the call sites do, and they differ: the ones
* before the rows are durable take the statement back, the one after them does not.
*
* <p>One lookup, because a point is consumed by the lookup that reads it: reading it twice with
* {@code execute=1} armed would have the first read spend the allowance and the point be gone before the
* second, which would silently leave the overwrite uncancelled.
*/
private void cancelTheOverwriteAt(String debugPointName, TableIf targetTable) {
if (!targetTable.getName().equals(DebugPointUtil.getDebugParamOrDefault(
debugPointName, "table_name", ""))) {
return;
}
LOG.info("debug point {} cancels the overwrite of {}", debugPointName, targetTable.getName());
cancel();
}
/**
* Publishes this overwrite by running its swap, with a last look at the cancellation flag taken under the
* lock the swap contends for.
*
* <p>{@link #run} reads the flag before the swap is issued, and the swap then waits for the table's write
* lock, so a cancellation that arrives during that wait is the one place a check before the swap cannot
* see. Reading it again here costs nothing and is where the wait happens: for a cancellation with nothing
* committed there is nothing durable to publish, so refusing to swap costs the statement and leaves the
* rows the client asked to keep -- while a swap that went ahead would replace them with an empty result.
*
* <p>Only a local table is wrapped: a remote table swaps on the frontend that owns it, where this lock
* says nothing.
*/
private void publishTheOverwrite(TableIf targetTable, InsertCommandContext insertCtx, ConnectContext ctx,
OverwritePublication publication) throws UserException {
if (!(targetTable instanceof OlapTable) || targetTable instanceof RemoteOlapTable) {
publication.publish();
return;
}
OlapTable olapTable = (OlapTable) targetTable;
if (!olapTable.writeLockIfExist()) {
// The target was dropped while this overwrite ran, so there is nothing to publish into and no swap
// to issue. Failing is also what the utility's own early return did for a dropped table -- its
// finally unlocks a lock that return never took, which raises -- and it is what a client whose
// swap never happened is owed: acknowledging the overwrite would claim rows the table cannot hold.
// The catch drops the temp partitions of the dropped table and takes the task back.
throw new UserException("insert overwrite could not publish its temporary partitions: table "
+ olapTable.getName() + " was dropped, queryId: " + ctx.getQueryIdentifier());
}
try {
cancelTheOverwriteAt(DEBUG_POINT_CANCEL_WHILE_THE_SWAP_WAITS_FOR_THE_TABLE_LOCK, targetTable);
if (isCancelled.get() && !insertCtx.hasCommitted()) {
throw cancelledBeforeTheRowsWereCommitted("while the swap waited for the table lock", ctx);
}
publication.publish();
} finally {
olapTable.writeUnlock();
}
}
/**
* The failure a cancellation that found nothing durable is reported as. No row and no stream offset was
* committed, so a re-run reads the same rows -- which is why this is a failure rather than the success of
* an overwrite that never ran. What a cancellation means on the other side of that boundary, where the
* rows are durable, is decided where the swap runs.
*/
private static UserException cancelledBeforeTheRowsWereCommitted(String stage, ConnectContext ctx) {
return new UserException("insert overwrite is cancelled " + stage + ", queryId: "
+ ctx.getQueryIdentifier());
}
private void runInsertCommand(LogicalPlan logicalQuery, InsertCommandContext insertCtx,
ConnectContext ctx, StmtExecutor executor) throws Exception {
InsertIntoTableCommand insertCommand = new InsertIntoTableCommand(logicalQuery, labelName,
Optional.of(insertCtx), Optional.empty(), false, branchName);
insertCommand.run(ctx, executor);
if (ctx.getState().getStateType() == MysqlStateType.ERR) {
if (insertCtx.hasCommitted()) {
// The rows are durable and only their publication timed out, which the session's
// visibility-timeout mode turns into this error (`insert_visible_timeout_return_mode=error`).
// Dropping the temp partitions for it would lose exactly what the error says was committed,
// so the overwrite keeps going: the swap below publishes the rows, and the error the client
// gets stays what it is -- a statement about visibility, not about whether the overwrite ran.
LOG.info("insert overwrite continues over an error state whose rows are committed, queryId: {}",
ctx.getQueryIdentifier());
return;
}
String errMsg = Strings.emptyToNull(ctx.getState().getErrorMessage());
LOG.warn("InsertInto state error:{}", errMsg);
throw new UserException(errMsg);
}
}
/**
* insert into select. for sepecified temp partitions or all partitions(table).
*
* @param ctx ctx
* @param executor executor
* @param tempPartitionNames tempPartitionNames
* @param wholeTable overwrite target is the whole table. not one by one by partitions(...)
* @return the context the inner insert ran under, which says whether it committed anything; see
* {@link InsertCommandContext#hasCommittedNothing()}
*/
private InsertCommandContext insertIntoPartitions(ConnectContext ctx, StmtExecutor executor,
List<String> tempPartitionNames, boolean wholeTable)
throws Exception {
// copy sink tot replace by tempPartitions
UnboundLogicalSink<?> copySink;
InsertCommandContext insertCtx;
LogicalPlan logicalQuery = getLogicalQuery();
if (logicalQuery instanceof UnboundTableSink) {
UnboundTableSink<?> sink = (UnboundTableSink<?>) logicalQuery;
copySink = (UnboundLogicalSink<?>) UnboundTableSinkCreator.createUnboundTableSink(
sink.getNameParts(),
sink.getColNames(),
sink.getHints(),
true,
tempPartitionNames,
sink.isPartialUpdate(),
sink.getPartialUpdateNewRowPolicy(),
sink.getDMLCommandType(),
(LogicalPlan) (sink.child(0)));
// 1. when overwrite table, allow auto partition or not is controlled by session variable.
// 2. we save and pass overwrite auto detect by insertCtx
boolean allowAutoPartition = wholeTable && ctx.getSessionVariable().isEnableAutoCreateWhenOverwrite();
insertCtx = new OlapInsertCommandContext(allowAutoPartition, true);
} else if (logicalQuery instanceof UnboundConnectorTableSink) {
UnboundConnectorTableSink<?> sink = (UnboundConnectorTableSink<?>) logicalQuery;
copySink = (UnboundLogicalSink<?>) UnboundTableSinkCreator.createUnboundTableSink(
sink.getNameParts(), sink.getColNames(), sink.getHints(),
false, sink.getPartitions(), false,
TPartialUpdateNewRowPolicy.APPEND,
sink.getDMLCommandType(),
(LogicalPlan) (sink.child(0)),
sink.getStaticPartitionKeyValues());
PluginDrivenInsertCommandContext pluginCtx = new PluginDrivenInsertCommandContext();
pluginCtx.setOverwrite(true);
// Thread the @branch target onto the generic write context (the inner InsertIntoTableCommand
// reuses this ctx) so the connector points the overwrite commit at the branch. The guard above
// already rejected @branch for connectors without supportsWriteBranch().
branchName.ifPresent(notUsed -> pluginCtx.setBranchName(branchName));
if (sink.hasStaticPartition()) {
Map<String, String> staticSpec = Maps.newHashMap();
for (Map.Entry<String, Expression> e : sink.getStaticPartitionKeyValues().entrySet()) {
if (e.getValue() instanceof Literal) {
staticSpec.put(e.getKey(), ((Literal) e.getValue()).getStringValue());
}
}
pluginCtx.setStaticPartitionSpec(staticSpec);
}
insertCtx = pluginCtx;
} else {
throw new UserException("Current catalog does not support insert overwrite yet.");
}
runInsertCommand(copySink, insertCtx, ctx, executor);
return insertCtx;
}
/**
* insert into auto detect partition.
*
* @param ctx ctx
* @param executor executor
* @return the context the inner insert ran under, which says whether it committed anything; see
* {@link InsertCommandContext#hasCommittedNothing()}
*/
private InsertCommandContext insertIntoAutoDetect(ConnectContext ctx, StmtExecutor executor, long groupId)
throws Exception {
InsertCommandContext insertCtx;
LogicalPlan logicalQuery = getLogicalQuery();
if (logicalQuery instanceof UnboundTableSink) {
// 1. when overwrite auto-detect, allow auto partition or not is controlled by session variable.
// 2. we save and pass overwrite auto detect by insertCtx
boolean allowAutoPartition = ctx.getSessionVariable().isEnableAutoCreateWhenOverwrite();
insertCtx = new OlapInsertCommandContext(allowAutoPartition,
((UnboundTableSink<?>) logicalQuery).isAutoDetectPartition(), groupId, true);
} else {
throw new UserException("Current catalog does not support insert overwrite with auto-detect partition.");
}
runInsertCommand(logicalQuery, insertCtx, ctx, executor);
return insertCtx;
}
@Override
public Plan getExplainPlan(ConnectContext ctx) {
Optional<CascadesContext> analyzeContext = Optional.of(
CascadesContext.initContext(ctx.getStatementContext(), originLogicalQuery, PhysicalProperties.ANY)
);
return InsertUtils.getPlanForExplain(ctx, analyzeContext, getLogicalQuery(), branchName);
}
@Override
public Optional<NereidsPlanner> getExplainPlanner(LogicalPlan logicalPlan, StatementContext ctx) {
LogicalPlan logicalQuery = getLogicalQuery();
if (logicalQuery instanceof UnboundTableSink) {
boolean allowAutoPartition = ctx.getConnectContext().getSessionVariable().isEnableAutoCreateWhenOverwrite();
OlapInsertCommandContext insertCtx = new OlapInsertCommandContext(allowAutoPartition, true);
InsertIntoTableCommand insertIntoTableCommand = new InsertIntoTableCommand(
logicalQuery, labelName, Optional.of(insertCtx), Optional.empty(), true, Optional.empty());
return insertIntoTableCommand.getExplainPlanner(logicalPlan, ctx);
}
return Optional.empty();
}
public boolean isForceDropPartition() {
return true;
}
@Override
public <R, C> R accept(PlanVisitor<R, C> visitor, C context) {
return visitor.visitInsertOverwriteTableCommand(this, context);
}
@Override
public StmtType stmtType() {
return StmtType.INSERT;
}
@Override
public boolean needAuditEncryption() {
return originLogicalQuery.anyMatch(node -> node instanceof TVFRelation);
}
@Override
public String toDigest() {
// if with cte, query will be print twice
StringBuilder sb = new StringBuilder();
sb.append("OVERWRITE TABLE "); // there is no way add overwrite flag in sink(logic query), so add it here
sb.append(originLogicalQuery.toDigest());
if (cte.isPresent()) {
sb.append(" (").append(cte.get().toDigest()).append(")");
}
return sb.toString();
}
}