github-actions[bot] commented on code in PR #68662:
URL: https://github.com/apache/doris/pull/68662#discussion_r4140574897
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertOverwriteTableCommand.java:
##########
@@ -386,6 +409,37 @@ private static void
failBetweenTheTwoHalvesOfAnOverwrite(TableIf targetTable) th
throw new UserException("debug point: " +
DEBUG_POINT_FAIL_BETWEEN_THE_HALVES_OF_AN_OVERWRITE);
}
+ /**
+ * Cancels this overwrite when the debug point names the stage and the
table it targets; see the constant
+ * above. Nothing here decides what a cancelled overwrite means -- the two
call sites do, and they differ:
+ * one takes the statement back, the other cannot.
+ */
+ private void cancelTheOverwriteAt(String stage, TableIf targetTable) {
+ if (!stage.equals(DebugPointUtil.getDebugParamOrDefault(
+ DEBUG_POINT_CANCEL_AN_OVERWRITE, "stage", ""))) {
Review Comment:
[P3] Make one-shot cancellation points usable at both stages.
`getDebugParamOrDefault` consumes the point's `execute` budget on every lookup.
With `execute=1` for `beforeTheInsert`, the `stage` read uses the one allowance
and this `table_name` read expires the point; for `afterTheInsert`, the earlier
nonmatching before-stage check consumes it. Both silently skip `cancel()`. Read
both parameters from one point instance and avoid consuming a stage's one-shot
budget at the other stage (for example, use stage-specific point names).
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertOverwriteTableCommand.java:
##########
@@ -266,32 +283,37 @@ public void run(ConnectContext ctx, StmtExecutor
executor) throws Exception {
} else {
// it's overwrite table(as all partitions) or specific
partition(s)
List<String> tempPartitionNames =
InsertOverwriteUtil.generateTempPartitionNames(partitionNames);
Review Comment:
[P2] Apply the pre-commit cancel rule to `PARTITION(*)` as well. This check
is only in the explicit-partition branch. For `INSERT OVERWRITE t PARTITION(*)
SELECT ... LIMIT 0`, the inner insert returns through its empty-plan fast path
without a transaction or coordinator, so a real KILL leaves the outer
cancellation flag set but the auto-detect branch still calls `taskGroupSuccess`
and reports success. Check cancellation in that branch when no insert
committed, so this form does not retain the false-success behavior the PR is
fixing.
##########
regression-test/suites/insert_overwrite_p0/test_insert_overwrite_cancel.groovy:
##########
@@ -0,0 +1,102 @@
+// 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.
+
+/**
+ * What a cancelled overwrite owes the rows it has already committed, and what
it owes the caller when it has
+ * not committed anything yet.
+ *
+ * <p>An overwrite commits its rows into temporary partitions and publishes
them with a swap afterwards, so a
+ * cancellation that lands between the two halves cannot take the rows back:
they are durable, and everything
+ * the write read -- the base table stream offsets among it -- was committed
with them. Dropping the temporary
+ * partitions there, and answering the client with a success for an overwrite
that never happened, is what
+ * loses those rows against an advanced offset. What the command has to do
instead is finish the swap. A
+ * cancellation that lands before the first half has committed anything has
nothing to take back, and there
+ * the statement has to fail rather than report the success of an overwrite
that did not run.
+ *
+ * <p>Both are injected by one debug point in the overwrite, scoped by stage
and by table name, so each case
+ * lands on a real statement rather than on a race.
+ */
+suite("test_insert_overwrite_cancel", "nonConcurrent") {
+ def cancelPoint = "InsertOverwriteTableCommand.cancelAnOverwrite"
+
+ GetDebugPoint().disableDebugPointForAllFEs(cancelPoint)
+ sql """DROP TABLE IF EXISTS test_iot_cancel_src"""
+ sql """DROP TABLE IF EXISTS test_iot_cancel_dst"""
+
+ sql """
+ CREATE TABLE test_iot_cancel_src (
+ id BIGINT NOT NULL,
+ dt DATE NOT NULL,
+ amount INT
+ ) ENGINE = OLAP
+ UNIQUE KEY(id, dt)
+ PARTITION BY RANGE(dt) (
+ PARTITION p1 VALUES [('2026-01-01'), ('2026-02-01')),
+ PARTITION p2 VALUES [('2026-02-01'), ('2026-03-01'))
+ )
+ DISTRIBUTED BY HASH(id) BUCKETS 1
+ PROPERTIES (
+ "replication_num" = "1",
+ "enable_unique_key_merge_on_write" = "true"
+ )
+ """
+ sql """
+ CREATE TABLE test_iot_cancel_dst (
+ id BIGINT NOT NULL,
+ dt DATE NOT NULL,
+ amount INT
+ ) ENGINE = OLAP
+ UNIQUE KEY(id, dt)
+ PARTITION BY RANGE(dt) (
+ PARTITION p1 VALUES [('2026-01-01'), ('2026-02-01')),
+ PARTITION p2 VALUES [('2026-02-01'), ('2026-03-01'))
+ )
+ DISTRIBUTED BY HASH(id) BUCKETS 1
+ PROPERTIES (
+ "replication_num" = "1",
+ "enable_unique_key_merge_on_write" = "true"
+ )
+ """
+ sql """INSERT INTO test_iot_cancel_dst VALUES (1, '2026-01-10', 100), (2,
'2026-02-10', 200)"""
+ sql """INSERT INTO test_iot_cancel_src VALUES (3, '2026-01-20', 300), (4,
'2026-02-20', 400)"""
+
+ // A cancellation that lands before the rows are committed has nothing to
take back, so the statement has
+ // to report it, and the table has to be where it was.
+ try {
+ GetDebugPoint().enableDebugPointForAllFEs(cancelPoint,
+ [stage: "beforeTheInsert", table_name: "test_iot_cancel_dst"])
+ test {
+ sql """INSERT OVERWRITE TABLE test_iot_cancel_dst SELECT * FROM
test_iot_cancel_src"""
+ exception "insert overwrite is cancelled before registerTask"
+ }
+ } finally {
+ GetDebugPoint().disableDebugPointForAllFEs(cancelPoint)
+ }
+ order_qt_dst_after_the_cancelled_statement """SELECT id, dt, amount FROM
test_iot_cancel_dst"""
+
+ // A cancellation that lands after the rows are committed cannot take them
back: they are in the temporary
+ // partitions, and the swap is what publishes them. The statement reports
the success it now is, and the
+ // rows it read are the ones the table holds.
+ try {
+ GetDebugPoint().enableDebugPointForAllFEs(cancelPoint,
Review Comment:
[P2] Make the post-insert cancellation case observable. If the
`afterTheInsert` call is removed or its stage stops matching, this SQL still
succeeds and the final query still returns rows 3 and 4, so the suite passes
without exercising the behavior it was added to protect. Assert that this stage
actually cancelled the command (and preferably exercise a real
`StmtExecutor.cancel`/stream-offset case) before treating the row result as a
regression guard.
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertOverwriteTableCommand.java:
##########
@@ -266,32 +283,37 @@ public void run(ConnectContext ctx, StmtExecutor
executor) throws Exception {
} else {
// it's overwrite table(as all partitions) or specific
partition(s)
List<String> tempPartitionNames =
InsertOverwriteUtil.generateTempPartitionNames(partitionNames);
+ cancelTheOverwriteAt(STAGE_BEFORE_THE_INSERT, targetTable);
if (isCancelled.get()) {
- LOG.info("insert overwrite is cancelled before
registerTask, queryId: {}",
- ctx.getQueryIdentifier());
- return;
+ // 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()) {
- LOG.info("insert overwrite is cancelled before
addTempPartitions, queryId: {}",
- ctx.getQueryIdentifier());
- // not need deal temp partition
- insertOverwriteManager.taskSuccess(taskId);
- return;
+ // 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()) {
- LOG.info("insert overwrite is cancelled before insertInto,
queryId: {}", ctx.getQueryIdentifier());
- insertOverwriteManager.taskFail(taskId);
- return;
+ // 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
insertIntoPartitions(ctx, executor, tempPartitionNames,
wholeTable);
+ cancelTheOverwriteAt(STAGE_AFTER_THE_INSERT, targetTable);
if (isCancelled.get()) {
- LOG.info("insert overwrite is cancelled before
replacePartition, queryId: {}",
- ctx.getQueryIdentifier());
- insertOverwriteManager.taskFail(taskId);
- return;
Review Comment:
[P1] Preserve cancellation when the inner insert has no transaction. `INSERT
OVERWRITE dst SELECT ... LIMIT 0` becomes a `PhysicalEmptyRelation`;
`InsertIntoTableCommand.runInternal` returns at `!requiresTransaction()`
without committing or publishing a coordinator. If KILL lands after the earlier
check (or this after-insert debug point fires), `isCancelled` is true here, yet
the command swaps the empty temporary partitions over `dst` and reports
success. This deletes the prior rows even though cancellation occurred before
any insert commit. Fail and roll back this no-transaction case; complete the
swap only after a committed insert.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]