github-actions[bot] commented on code in PR #68662:
URL: https://github.com/apache/doris/pull/68662#discussion_r4142311115


##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/RemoteOlapInsertExecutor.java:
##########
@@ -208,6 +208,7 @@ protected void onComplete() throws UserException {
                 } else {
                     txnStatus = TransactionStatus.COMMITTED;
                 }
+                markCommitted();
                 LOG.info("commit remote txn success, catalog={}, dbId={}, 
txnId={}, status={}",

Review Comment:
   [P1] Reconcile a remote commit whose reply is lost before rolling back the 
overwrite. The owning FE commits before it waits for publication, while 
`masterCallWithRetry` can exhaust retries after a lost reply. In that case this 
marker stays false although the remote transaction is COMMITTED; `onFail` 
cannot abort it, and the overwrite treats its ERR state as uncommitted and 
calls `taskFail`, dropping the partitions containing those rows. Query the 
remote transaction outcome or preserve its temporary partitions while the 
result is unknown. This differs from the existing visibility-timeout thread, 
where the committed response reached this FE.



##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/OlapInsertExecutor.java:
##########
@@ -240,10 +240,14 @@ protected void onComplete() throws UserException {
                 ctx.getSessionVariable().getInsertVisibleTimeoutMs(), 
txnCommitAttachment,
                 streamUpdateInfos)) {
             txnStatus = TransactionStatus.VISIBLE;
+            markCommitted();
         } else {
             // Keep the committed status so load accounting and insert result 
bookkeeping stay aligned.
             txnStatus = TransactionStatus.COMMITTED;
             publishTimedOutAfterCommit = true;
+            // Committed, visible later: the rows are durable even though the 
session's visibility-timeout
+            // mode may report the timeout as an error. See 
InsertCommandContext#setCommitted.
+            markCommitted();
         }

Review Comment:
   [P1] Resolve an ambiguous cloud commit before failing the overwrite. 
FoundationDB can report a maybe-committed result after applying the 
transaction; after finite meta-service retries this reaches FE as 
`KV_TXN_COMMIT_ERR`. `commitAndPublishTransactionWithRetry` then throws before 
either `markCommitted()` call, leaving this context false. The overwrite sees 
ERR and `taskFail` drops its temp partitions even if the cloud commit made 
their rows and stream offsets durable. Determine the cloud transaction state, 
or keep the partitions when commit remains uncertain, before permitting 
rollback. This is separate from the returned visibility timeout already 
discussed.



##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertOverwriteTableCommand.java:
##########
@@ -386,12 +440,87 @@ 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 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 swapping the temp partitions in, 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 empty partitions.
+     *
+     * <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, List<String> 
partitionNames,

Review Comment:
   [P1] Carry precommit cancellation through the remote swap. For an 
unpartitioned remote target, `SELECT ... WHERE 1 = 0` can return from the inner 
insert without a transaction. If KILL arrives after the check at line 332 while 
the owning FE waits for its table write lock in `replacePartitionsImpl`, this 
branch supplies no cancellation state or last check: the empty temp partition 
replaces existing rows and the statement reports success. The existing local 
lock-wait fix does not cover this remote path; make the owning FE decide under 
its swap lock before replacing an uncommitted empty result.



##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertOverwriteTableCommand.java:
##########
@@ -386,12 +440,87 @@ 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 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 swapping the temp partitions in, 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 empty partitions.
+     *
+     * <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, List<String> 
partitionNames,
+            List<String> tempPartitionNames, InsertCommandContext insertCtx, 
ConnectContext ctx) throws UserException {
+        if (!(targetTable instanceof OlapTable) || targetTable instanceof 
RemoteOlapTable) {
+            InsertOverwriteUtil.replacePartition(targetTable, partitionNames, 
tempPartitionNames,
+                    isForceDropPartition());
+            return;
+        }

Review Comment:
   [P2] Fail when the target was dropped before the swap. `writeLockIfExist()` 
returns false after a concurrent DROP, so this new early return skips 
`replacePartition`, but `run` still calls `taskSuccess` and reports the 
overwrite as successful. The previous direct utility path raised an exception 
in this case; preserve an error outcome and clean the task rather than 
acknowledging a swap that never occurred.



-- 
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