This is an automated email from the ASF dual-hosted git repository.

mymeiyi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/master by this push:
     new 72071af801d [fix](fe) Return transaction insert execution errors 
(#67386)
72071af801d is described below

commit 72071af801d265f87511bb1072cc9c6aeba2d1d2
Author: meiyi <[email protected]>
AuthorDate: Wed Sep 2 18:41:20 2026 +0800

    [fix](fe) Return transaction insert execution errors (#67386)
    
    Transaction INSERT execution failures raised directly while dispatching 
fragments bypass ErrorReport. OlapTxnInsertExecutor catches these failures and 
aborts the subtransaction, but previously left QueryState successful, causing 
the client to receive Query OK. Set the INSERT error state when no earlier 
error exists so the failure reaches the client while preserving errors already 
reported by ErrorReport.
---
 .../transaction/CloudGlobalTransactionMgr.java     |  2 +-
 .../insert/OlapGroupCommitInsertExecutor.java      | 17 +----------
 .../plans/commands/insert/OlapInsertExecutor.java  |  6 +++-
 .../commands/insert/OlapTxnInsertExecutor.java     |  1 +
 .../commands/insert/RemoteOlapInsertExecutor.java  | 19 +-----------
 .../transaction/txn_insert_inject_case.groovy      | 35 +++++++++++++++-------
 6 files changed, 34 insertions(+), 46 deletions(-)

diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgr.java
 
b/fe/fe-core/src/main/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgr.java
index d0b97593cec..fa2fd8bcf08 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgr.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgr.java
@@ -2777,7 +2777,7 @@ public class CloudGlobalTransactionMgr implements 
GlobalTransactionMgrIface {
         LOG.info("try to abort sub transaction, txnId: {}, subTxnId: {}, dbId: 
{}, tableIds: {}, subTxnNum: {}", txnId,
                 subTxnId, dbId, tableIds, subTxnNum);
         AbortSubTxnRequest request = 
AbortSubTxnRequest.newBuilder().setCloudUniqueId(Config.cloud_unique_id)
-                
.setTxnId(txnId).setSubTxnId(subTxnId).setDbId(dbId).addAllTableIds(tableIds).setSubTxnNum(subTxnId)
+                
.setTxnId(txnId).setSubTxnId(subTxnId).setDbId(dbId).addAllTableIds(tableIds).setSubTxnNum(subTxnNum)
                 
.setRequestIp(FrontendOptions.getLocalHostAddressCached()).build();
         AbortSubTxnResponse response = null;
         int retryTime = 0;
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/OlapGroupCommitInsertExecutor.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/OlapGroupCommitInsertExecutor.java
index 0b4d9026d91..5b2fe260b3b 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/OlapGroupCommitInsertExecutor.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/OlapGroupCommitInsertExecutor.java
@@ -22,8 +22,6 @@ import org.apache.doris.catalog.MTMV;
 import org.apache.doris.catalog.OlapTable;
 import org.apache.doris.catalog.Table;
 import org.apache.doris.catalog.TableIf;
-import org.apache.doris.common.Config;
-import org.apache.doris.common.ErrorCode;
 import org.apache.doris.common.FeConstants;
 import org.apache.doris.common.Pair;
 import org.apache.doris.common.util.DebugUtil;
@@ -44,8 +42,6 @@ import org.apache.doris.qe.StmtExecutor;
 import org.apache.doris.system.Backend;
 import org.apache.doris.transaction.TransactionStatus;
 
-import com.google.common.base.Strings;
-import org.apache.commons.lang3.StringUtils;
 import org.apache.logging.log4j.LogManager;
 import org.apache.logging.log4j.Logger;
 
@@ -192,18 +188,7 @@ public class OlapGroupCommitInsertExecutor extends 
OlapInsertExecutor {
         String queryId = DebugUtil.printId(ctx.queryId());
         // if any throwable being thrown during insert operation, first we 
should abort this txn
         LOG.warn("insert [{}] with query id {} failed, url={}", labelName, 
queryId, coordinator.getTrackingUrl(), t);
-        String firstErrorMsgPart = "";
-        String urlPart = "";
-        if (!Strings.isNullOrEmpty(coordinator.getFirstErrorMsg())) {
-            firstErrorMsgPart = 
StringUtils.abbreviate(coordinator.getFirstErrorMsg(),
-                    Config.first_error_msg_max_length);
-        }
-        if (!Strings.isNullOrEmpty(coordinator.getTrackingUrl())) {
-            urlPart = coordinator.getTrackingUrl();
-        }
-
-        String finalErrorMsg = InsertUtils.getFinalErrorMsg(errMsg, 
firstErrorMsgPart, urlPart);
-        ctx.getState().setError(ErrorCode.ERR_UNKNOWN_ERROR, finalErrorMsg);
+        setErrorState();
     }
 
     @Override
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/OlapInsertExecutor.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/OlapInsertExecutor.java
index 2563077a30f..7dc65ebce29 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/OlapInsertExecutor.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/OlapInsertExecutor.java
@@ -312,6 +312,11 @@ public class OlapInsertExecutor extends 
AbstractInsertExecutor {
         if (Config.isCloudMode() && 
SystemInfoService.needRetryWithReplan(t.getMessage())) {
             return;
         }
+        setErrorState();
+        recordLoadJob(ctx.getCurrentUserIdentity());
+    }
+
+    protected void setErrorState() {
         String firstErrorMsgPart = "";
         String urlPart = "";
         if (!Strings.isNullOrEmpty(coordinator.getFirstErrorMsg())) {
@@ -323,7 +328,6 @@ public class OlapInsertExecutor extends 
AbstractInsertExecutor {
         }
         String finalErrorMsg = InsertUtils.getFinalErrorMsg(errMsg, 
firstErrorMsgPart, urlPart);
         ctx.getState().setError(ErrorCode.ERR_UNKNOWN_ERROR, finalErrorMsg);
-        recordLoadJob(ctx.getCurrentUserIdentity());
     }
 
     private void recordLoadJob(UserIdentity userIdentity) {
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/OlapTxnInsertExecutor.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/OlapTxnInsertExecutor.java
index aa095d3f91a..ad7ee6c067b 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/OlapTxnInsertExecutor.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/OlapTxnInsertExecutor.java
@@ -87,6 +87,7 @@ public class OlapTxnInsertExecutor extends OlapInsertExecutor 
{
         // if any throwable being thrown during insert operation, first we 
should abort this txn
         LOG.warn("insert [{}] with query id {} failed, url={}", labelName, 
queryId, coordinator.getTrackingUrl(), t);
         cleanTransaction();
+        setErrorState();
     }
 
     @Override
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/RemoteOlapInsertExecutor.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/RemoteOlapInsertExecutor.java
index 0003d44f633..c95c3b0d71d 100755
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/RemoteOlapInsertExecutor.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/RemoteOlapInsertExecutor.java
@@ -19,7 +19,6 @@ package org.apache.doris.nereids.trees.plans.commands.insert;
 
 import org.apache.doris.catalog.Env;
 import org.apache.doris.common.AuthenticationException;
-import org.apache.doris.common.ErrorCode;
 import org.apache.doris.common.LabelAlreadyUsedException;
 import org.apache.doris.common.MetaNotFoundException;
 import org.apache.doris.common.QuotaExceedException;
@@ -54,7 +53,6 @@ import org.apache.doris.transaction.BeginTransactionException;
 import org.apache.doris.transaction.TransactionStatus;
 
 import com.google.common.base.Strings;
-import org.apache.commons.lang3.StringUtils;
 import org.apache.logging.log4j.LogManager;
 import org.apache.logging.log4j.Logger;
 
@@ -260,20 +258,6 @@ public class RemoteOlapInsertExecutor extends 
OlapInsertExecutor {
         }
     }
 
-    private String buildFinalErrorMessage(Throwable t) {
-        String localErrMsg = t.getMessage() == null ? "unknown reason" : 
t.getMessage();
-        String firstErrorMsgPart = "";
-        String urlPart = "";
-        if (!Strings.isNullOrEmpty(coordinator.getFirstErrorMsg())) {
-            firstErrorMsgPart = 
StringUtils.abbreviate(coordinator.getFirstErrorMsg(),
-                    org.apache.doris.common.Config.first_error_msg_max_length);
-        }
-        if (!Strings.isNullOrEmpty(coordinator.getTrackingUrl())) {
-            urlPart = coordinator.getTrackingUrl();
-        }
-        return InsertUtils.getFinalErrorMsg(localErrMsg, firstErrorMsgPart, 
urlPart);
-    }
-
     @Override
     protected void onFail(Throwable t) {
         errMsg = t.getMessage() == null ? "unknown reason" : t.getMessage();
@@ -287,8 +271,7 @@ public class RemoteOlapInsertExecutor extends 
OlapInsertExecutor {
                         labelName, queryId, txnId, abortTxnException);
             }
         }
-        String finalErrorMsg = buildFinalErrorMessage(t);
-        ctx.getState().setError(ErrorCode.ERR_UNKNOWN_ERROR, finalErrorMsg);
+        setErrorState();
     }
 
     @Override
diff --git 
a/regression-test/suites/insert_p0/transaction/txn_insert_inject_case.groovy 
b/regression-test/suites/insert_p0/transaction/txn_insert_inject_case.groovy
index 65556a91301..39620e8d4b1 100644
--- a/regression-test/suites/insert_p0/transaction/txn_insert_inject_case.groovy
+++ b/regression-test/suites/insert_p0/transaction/txn_insert_inject_case.groovy
@@ -40,6 +40,7 @@ suite("txn_insert_inject_case", "nonConcurrent") {
         """
     }
     GetDebugPoint().disableDebugPointForAllBEs("FlushToken.submit_flush_error")
+    
GetDebugPoint().disableDebugPointForAllBEs("FragmentMgr.exec_plan_fragment.failed")
     sql """insert into ${table}_1 values(1, 2.2, "abc", [], []), (2, 3.3, 
"xyz", [1], [1, 0]), (null, null, null, [null], [null, 0])  """
     sql """insert into ${table}_2 values(3, 2.2, "abc", [], []), (4, 3.3, 
"xyz", [1], [1, 0]), (null, null, null, [null], [null, 0])  """
 
@@ -48,20 +49,20 @@ suite("txn_insert_inject_case", "nonConcurrent") {
     (ipList, portList) = GetDebugPoint().getBEHostAndHTTPPort()
     logger.info("be ips: ${ipList}, ports: ${portList}")
 
-    def enableDebugPoint = { ->
+    def enableDebugPoint = { debugPoint ->
         ipList.each { beid, ip ->
-            DebugPoint.enableDebugPoint(ip, portList[beid] as int, 
NodeType.BE, "FlushToken.submit_flush_error")
+            DebugPoint.enableDebugPoint(ip, portList[beid] as int, 
NodeType.BE, debugPoint)
         }
     }
 
-    def disableDebugPoint = { ->
+    def disableDebugPoint = { debugPoint ->
         ipList.each { beid, ip ->
-            DebugPoint.disableDebugPoint(ip, portList[beid] as int, 
NodeType.BE, "FlushToken.submit_flush_error")
+            DebugPoint.disableDebugPoint(ip, portList[beid] as int, 
NodeType.BE, debugPoint)
         }
     }
 
     try {
-        enableDebugPoint()
+        enableDebugPoint("FlushToken.submit_flush_error")
         sql """ begin """
         try {
             sql """ insert into ${table}_0 select * from ${table}_1; """
@@ -78,10 +79,10 @@ suite("txn_insert_inject_case", "nonConcurrent") {
             
assertTrue(e.getMessage().contains("dbug_be_memtable_submit_flush_error"))
         }
 
-        disableDebugPoint()
+        disableDebugPoint("FlushToken.submit_flush_error")
         sql """ insert into ${table}_0 select * from ${table}_1; """
 
-        enableDebugPoint()
+        enableDebugPoint("FlushToken.submit_flush_error")
         try {
             sql """ insert into ${table}_0 select * from ${table}_1; """
             assertTrue(false, "insert should fail")
@@ -90,16 +91,30 @@ suite("txn_insert_inject_case", "nonConcurrent") {
             
assertTrue(e.getMessage().contains("dbug_be_memtable_submit_flush_error"))
         }
 
-        disableDebugPoint()
-        sql """ insert into ${table}_0 select * from ${table}_1; """
+        disableDebugPoint("FlushToken.submit_flush_error")
         sql """ commit"""
     } catch (Exception e) {
         logger.error("failed", e)
     } finally {
         sql """ rollback """
-        disableDebugPoint()
+        disableDebugPoint("FlushToken.submit_flush_error")
         
GetDebugPoint().disableDebugPointForAllBEs("FlushToken.submit_flush_error")
     }
+
+    try {
+        enableDebugPoint("FragmentMgr.exec_plan_fragment.failed")
+        sql """ begin """
+        test {
+            sql """ insert into ${table}_0 select * from ${table}_1; """
+            exception "FragmentMgr.exec_plan_fragment.failed"
+        }
+        disableDebugPoint("FragmentMgr.exec_plan_fragment.failed")
+        sql """ insert into ${table}_0 select * from ${table}_1; """
+        sql """ commit """
+    } finally {
+        sql """ rollback """
+        disableDebugPoint("FragmentMgr.exec_plan_fragment.failed")
+    }
     sql "sync"
     order_qt_select1 """select * from ${table}_0"""
 


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to