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]

Reply via email to