This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new ffb06fdf3c0 branch-4.1: [fix](fe) Return transaction insert execution
errors (#67386) (#67925)
ffb06fdf3c0 is described below
commit ffb06fdf3c0de2241b1a18dd09c7609eeb5d39b3
Author: meiyi <[email protected]>
AuthorDate: Tue Sep 15 11:48:18 2026 +0800
branch-4.1: [fix](fe) Return transaction insert execution errors (#67386)
(#67925)
pick https://github.com/apache/doris/pull/67386
---
.../transaction/CloudGlobalTransactionMgr.java | 2 +-
.../insert/OlapGroupCommitInsertExecutor.java | 17 +----------
.../plans/commands/insert/OlapInsertExecutor.java | 6 +++-
.../commands/insert/OlapTxnInsertExecutor.java | 1 +
.../transaction/txn_insert_inject_case.groovy | 35 +++++++++++++++-------
5 files changed, 33 insertions(+), 28 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 df0e193b0e4..1c49b3c11b2 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
@@ -2681,7 +2681,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 3771c415c17..6cbc9d40c88 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
@@ -307,6 +307,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())) {
@@ -318,7 +323,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/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]