yujun777 commented on code in PR #68662:
URL: https://github.com/apache/doris/pull/68662#discussion_r4140909403
##########
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:
Applied the same boundary to this branch: `insertIntoAutoDetect` now hands
its context back, and a cancellation there is honoured when the load committed
nothing (`hasCommittedNothing()`), the same way the explicit-partition branch
does.
One correction to the premise, from probing the plans on a local cluster: a
`PARTITION(*)` target is partitioned, and an empty query over a partitioned
target plans an exchange, so the sink's child is not the
`PhysicalEmptyRelation` and `requiresTransaction()` is true -- the load does
begin and commit an (empty) transaction rather than returning through the
no-transaction path. The empty-plan shape reproduces on an unpartitioned
target, which is what the new suite case uses (`INSERT OVERWRITE TABLE flat_dst
SELECT ... WHERE 1 = 0`). The branch is kept so both routes take the same
decision.
##########
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:
Fixed. The insert now tells its caller that nothing was committed:
`InsertIntoTableCommand.runInternal` marks the context at the
`!requiresTransaction()` return (`InsertCommandContext#setCommittedNothing`),
the overwrite reads that context back from `insertIntoPartitions`, and the
window check fails the statement -- rolling the empty temporary partitions back
through the existing catch -- instead of completing the swap.
Verified end to end on an unpartitioned empty-plan overwrite: with the
cancellation injected in the window, `INSERT OVERWRITE TABLE dst SELECT ...
WHERE 1 = 0` now returns `insert overwrite is cancelled after an insert that
committed nothing` and the target keeps its rows. Before the change the
statement succeeded and the table was emptied (FE log `cancelled after its rows
were committed, completing it`, `count(*)` 1 -> 0). The new suite case pins it.
##########
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:
Fixed: two point names instead of one, and one lookup per site.
`cancelBeforeTheInsertOfAnOverwrite` and
`cancelBetweenTheTwoHalvesOfAnOverwrite` each read their single `table_name`
parameter from one `getDebugPoint` call, so no site can spend the other's
`execute` allowance; the `stage` parameter and the second lookup are gone. The
comment on the constants records why a shared name was wrong.
##########
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:
Fixed by making the window's injection observable, in the way the earlier
failure-half suite does it. Alongside the case that asserts the rows are
published, the same injection point now drives an empty-plan overwrite (`...
WHERE 1 = 0`) where the cancelled statement has to fail: that assertion holds
only if the injection reached the window, so deleting the call, renaming the
point, or breaking its lookup makes the suite fail. The empty-plan case also
pins its own behaviour (the table has to keep its rows, so the swap must not
have run) rather than passing on the rows either way.
A real `StmtExecutor.cancel` is not used for the same reason the failure
half is injected: the window is a couple of metadata operations wide, so a KILL
would have to land inside it by luck. The stream-offset variant of this
scenario is the one that matters for IVM refreshes; it needs the crash half of
the window (swap as a committed action of the insert transaction) to be worth
testing end to end.
--
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]