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]