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

morningman pushed a commit to branch branch-4.0
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/branch-4.0 by this push:
     new 7d1d5eba3b0 branch-4.0: [feat](job) show first error message for 
routine load (#65661)
7d1d5eba3b0 is described below

commit 7d1d5eba3b01a5dc52449e1ea473cbe242204f92
Author: hui lai <[email protected]>
AuthorDate: Fri Jul 17 12:59:08 2026 +0800

    branch-4.0: [feat](job) show first error message for routine load (#65661)
    
    pick apache/doris#65553
---
 .../schema_routine_load_job_scanner.cpp            |   4 +
 .../runtime/stream_load/stream_load_executor.cpp   |   7 +-
 .../java/org/apache/doris/catalog/SchemaTable.java |   1 +
 .../transaction/CloudGlobalTransactionMgr.java     |   6 +-
 .../apache/doris/cloud/transaction/TxnUtil.java    |   3 +
 .../routineload/RLTaskTxnCommitAttachment.java     |  15 +++
 .../doris/load/routineload/RoutineLoadJob.java     |  35 ++++--
 .../plans/commands/ShowRoutineLoadCommand.java     |   1 +
 .../apache/doris/service/FrontendServiceImpl.java  |   1 +
 .../transaction/CloudGlobalTransactionMgrTest.java |  52 +++++++++
 .../doris/load/routineload/RoutineLoadJobTest.java |  39 +++++++
 .../plans/commands/ShowRoutineLoadCommandTest.java |   2 +
 gensrc/proto/cloud.proto                           |   2 +
 gensrc/thrift/FrontendService.thrift               |   2 +
 .../test_routine_load_first_error_msg.groovy       | 127 +++++++++++++++++++++
 15 files changed, 284 insertions(+), 13 deletions(-)

diff --git a/be/src/exec/schema_scanner/schema_routine_load_job_scanner.cpp 
b/be/src/exec/schema_scanner/schema_routine_load_job_scanner.cpp
index c965338d46a..c661373d4c0 100644
--- a/be/src/exec/schema_scanner/schema_routine_load_job_scanner.cpp
+++ b/be/src/exec/schema_scanner/schema_routine_load_job_scanner.cpp
@@ -57,6 +57,7 @@ std::vector<SchemaScanner::ColumnDesc> 
SchemaRoutineLoadJobScanner::_s_tbls_colu
         {"CURRENT_ABORT_TASK_NUM", TYPE_INT, sizeof(int32_t), true},
         {"IS_ABNORMAL_PAUSE", TYPE_BOOLEAN, sizeof(int8_t), true},
         {"COMPUTE_GROUP", TYPE_STRING, sizeof(StringRef), true},
+        {"FIRST_ERROR_MSG", TYPE_STRING, sizeof(StringRef), true},
 };
 
 SchemaRoutineLoadJobScanner::SchemaRoutineLoadJobScanner()
@@ -177,6 +178,9 @@ Status 
SchemaRoutineLoadJobScanner::_fill_block_impl(vectorized::Block* block) {
                 case 20: // COMPUTE_GROUP
                     column_value = job_info.__isset.compute_group ? 
job_info.compute_group : "";
                     break;
+                case 21: // FIRST_ERROR_MSG
+                    column_value = job_info.__isset.first_error_msg ? 
job_info.first_error_msg : "";
+                    break;
                 }
 
                 str_refs[row_idx] =
diff --git a/be/src/runtime/stream_load/stream_load_executor.cpp 
b/be/src/runtime/stream_load/stream_load_executor.cpp
index 907fbdf7410..51bbf40edeb 100644
--- a/be/src/runtime/stream_load/stream_load_executor.cpp
+++ b/be/src/runtime/stream_load/stream_load_executor.cpp
@@ -99,6 +99,9 @@ Status 
StreamLoadExecutor::execute_plan_fragment(std::shared_ptr<StreamLoadConte
                 *status = Status::DataQualityError("too many filtered rows");
             }
         }
+        if (ctx->load_type == TLoadType::ROUTINE_LOAD || !status->ok()) {
+            ctx->first_error_msg = state->get_first_error_msg();
+        }
 
         if (status->ok()) {
             
DorisMetrics::instance()->stream_receive_bytes_total->increment(ctx->receive_bytes);
@@ -107,7 +110,6 @@ Status 
StreamLoadExecutor::execute_plan_fragment(std::shared_ptr<StreamLoadConte
             LOG(WARNING) << "fragment execute failed"
                          << ", err_msg=" << status->to_string() << ", " << 
ctx->brief();
             ctx->number_loaded_rows = 0;
-            ctx->first_error_msg = state->get_first_error_msg();
             // cancel body_sink, make sender known it
             if (ctx->body_sink != nullptr) {
                 ctx->body_sink->cancel(status->to_string());
@@ -396,6 +398,9 @@ bool 
StreamLoadExecutor::collect_load_stat(StreamLoadContext* ctx, TTxnCommitAtt
         rl_attach.__set_receivedBytes(ctx->receive_bytes);
         rl_attach.__set_loadedBytes(ctx->loaded_bytes);
         rl_attach.__set_loadCostMs(ctx->load_cost_millis);
+        if (!ctx->first_error_msg.empty()) {
+            rl_attach.__set_firstErrorMsg(ctx->first_error_msg);
+        }
 
         attach->rlTaskTxnCommitAttachment = rl_attach;
         attach->__isset.rlTaskTxnCommitAttachment = true;
diff --git a/fe/fe-core/src/main/java/org/apache/doris/catalog/SchemaTable.java 
b/fe/fe-core/src/main/java/org/apache/doris/catalog/SchemaTable.java
index 11d1ec3a8f3..9306dbc3d27 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/catalog/SchemaTable.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/catalog/SchemaTable.java
@@ -673,6 +673,7 @@ public class SchemaTable extends Table {
                                     .column("CURRENT_ABORT_TASK_NUM", 
ScalarType.createType(PrimitiveType.INT))
                                     .column("IS_ABNORMAL_PAUSE", 
ScalarType.createType(PrimitiveType.BOOLEAN))
                                     .column("COMPUTE_GROUP", 
ScalarType.createStringType())
+                                    .column("FIRST_ERROR_MSG", 
ScalarType.createStringType())
                                     .build())
             )
             .put("load_jobs",
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 5ac29a25c8e..839060215a6 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
@@ -1826,7 +1826,7 @@ public class CloudGlobalTransactionMgr implements 
GlobalTransactionMgrIface {
 
         AbortTxnResponse abortTxnResponse = null;
         try {
-            abortTxnResponse = abortTransactionImpl(dbId, transactionId, 
reason, null);
+            abortTxnResponse = abortTransactionImpl(dbId, transactionId, 
reason, txnCommitAttachment);
         } finally {
             handleAfterAbort(abortTxnResponse, txnCommitAttachment, 
transactionId,
                     callbackInfo.first, callbackInfo.second);
@@ -1890,6 +1890,10 @@ public class CloudGlobalTransactionMgr implements 
GlobalTransactionMgrIface {
         builder.setTxnId(transactionId);
         builder.setReason(reason);
         builder.setCloudUniqueId(Config.cloud_unique_id);
+        if (txnCommitAttachment instanceof RLTaskTxnCommitAttachment) {
+            builder.setCommitAttachment(TxnUtil.rlTaskTxnCommitAttachmentToPb(
+                    (RLTaskTxnCommitAttachment) txnCommitAttachment));
+        }
 
         final AbortTxnRequest abortTxnRequest = builder.build();
         AbortTxnResponse abortTxnResponse = null;
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/cloud/transaction/TxnUtil.java 
b/fe/fe-core/src/main/java/org/apache/doris/cloud/transaction/TxnUtil.java
index b947ef0649e..b557fa0ecac 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/cloud/transaction/TxnUtil.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/cloud/transaction/TxnUtil.java
@@ -256,6 +256,9 @@ public class TxnUtil {
         if (rtTaskTxnCommitAttachment.getErrorLogUrl() != null) {
             builder.setErrorLogUrl(rtTaskTxnCommitAttachment.getErrorLogUrl());
         }
+        if (rtTaskTxnCommitAttachment.getFirstErrorMsg() != null) {
+            
builder.setFirstErrorMsg(rtTaskTxnCommitAttachment.getFirstErrorMsg());
+        }
 
         attachementBuilder.setRlTaskTxnCommitAttachment(builder.build());
         return attachementBuilder.build();
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/RLTaskTxnCommitAttachment.java
 
b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/RLTaskTxnCommitAttachment.java
index 7b4f681cc42..c1a1c7904ae 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/RLTaskTxnCommitAttachment.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/RLTaskTxnCommitAttachment.java
@@ -18,12 +18,14 @@
 package org.apache.doris.load.routineload;
 
 import org.apache.doris.cloud.proto.Cloud.RLTaskTxnCommitAttachmentPB;
+import org.apache.doris.common.Config;
 import org.apache.doris.thrift.TRLTaskTxnCommitAttachment;
 import org.apache.doris.thrift.TUniqueId;
 import org.apache.doris.transaction.TransactionState;
 import org.apache.doris.transaction.TxnCommitAttachment;
 
 import com.google.gson.annotations.SerializedName;
+import org.apache.commons.lang3.StringUtils;
 
 // {"progress": "", "backendId": "", "taskSignature": "", "numOfErrorData": "",
 // "numOfTotalData": "", "taskId": "", "jobId": ""}
@@ -44,6 +46,7 @@ public class RLTaskTxnCommitAttachment extends 
TxnCommitAttachment {
     @SerializedName(value = "pro")
     private RoutineLoadProgress progress;
     private String errorLogUrl;
+    private String firstErrorMsg;
 
     public RLTaskTxnCommitAttachment() {
         super(TransactionState.LoadJobSourceType.ROUTINE_LOAD_TASK);
@@ -70,6 +73,9 @@ public class RLTaskTxnCommitAttachment extends 
TxnCommitAttachment {
         if (rlTaskTxnCommitAttachment.isSetErrorLogUrl()) {
             this.errorLogUrl = rlTaskTxnCommitAttachment.getErrorLogUrl();
         }
+        if (rlTaskTxnCommitAttachment.isSetFirstErrorMsg()) {
+            this.firstErrorMsg = 
abbreviateFirstErrorMsg(rlTaskTxnCommitAttachment.getFirstErrorMsg());
+        }
     }
 
     public RLTaskTxnCommitAttachment(RLTaskTxnCommitAttachmentPB 
rlTaskTxnCommitAttachment) {
@@ -87,6 +93,11 @@ public class RLTaskTxnCommitAttachment extends 
TxnCommitAttachment {
 
         this.progress = progress;
         this.errorLogUrl = rlTaskTxnCommitAttachment.getErrorLogUrl();
+        this.firstErrorMsg = 
abbreviateFirstErrorMsg(rlTaskTxnCommitAttachment.getFirstErrorMsg());
+    }
+
+    private static String abbreviateFirstErrorMsg(String firstErrorMsg) {
+        return StringUtils.abbreviate(firstErrorMsg, 
Config.first_error_msg_max_length);
     }
 
     public long getJobId() {
@@ -129,6 +140,10 @@ public class RLTaskTxnCommitAttachment extends 
TxnCommitAttachment {
         return errorLogUrl;
     }
 
+    public String getFirstErrorMsg() {
+        return firstErrorMsg;
+    }
+
     @Override
     public String toString() {
         return "RLTaskTxnCommitAttachment [filteredRows=" + filteredRows
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/RoutineLoadJob.java
 
b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/RoutineLoadJob.java
index 37c9a6a1c52..a259dca557c 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/RoutineLoadJob.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/RoutineLoadJob.java
@@ -272,8 +272,10 @@ public abstract class RoutineLoadJob
     protected Expr deleteCondition;
     // TODO(ml): error sample
 
-    // save the latest 3 error log urls
-    private Queue<String> errorLogUrls = EvictingQueue.create(3);
+    // Save the latest 3 error log URLs in memory. The corresponding first 
error message
+    // uses the same lifecycle and should not be persisted with the job.
+    private transient Queue<String> errorLogUrls = EvictingQueue.create(3);
+    private transient String firstErrorMsg = "";
 
     @SerializedName("ccid")
     private String cloudClusterId;
@@ -809,6 +811,10 @@ public abstract class RoutineLoadJob
         return errorLogUrls;
     }
 
+    public String getFirstErrorMsg() {
+        return Strings.nullToEmpty(firstErrorMsg);
+    }
+
     // RoutineLoadScheduler will run this method at fixed interval, and renew 
the timeout tasks
     public void processTimeoutTasks() {
         writeLock();
@@ -954,10 +960,7 @@ public abstract class RoutineLoadJob
                                 + "when current total rows is more than base 
or the filter ratio is more than the max")
                         .build());
             }
-            // reset currentTotalNum, currentErrorNum and otherMsg
-            this.jobStatistic.currentErrorRows = 0;
-            this.jobStatistic.currentTotalRows = 0;
-            this.otherMsg = "";
+            resetCurrentErrorStatistics();
             this.jobStatistic.currentAbortedTaskNum = 0;
         } else if (this.jobStatistic.currentErrorRows > maxErrorNum
                 || (this.jobStatistic.currentTotalRows > 0
@@ -977,13 +980,18 @@ public abstract class RoutineLoadJob
                         "current error rows is more than max_error_number "
                             + "or the max_filter_ratio is more than the value 
set"), isReplay);
             }
-            // reset currentTotalNum, currentErrorNum and otherMsg
-            this.jobStatistic.currentErrorRows = 0;
-            this.jobStatistic.currentTotalRows = 0;
-            this.otherMsg = "";
+            resetCurrentErrorStatistics();
         }
     }
 
+    private void resetCurrentErrorStatistics() {
+        this.jobStatistic.currentErrorRows = 0;
+        this.jobStatistic.currentTotalRows = 0;
+        this.otherMsg = "";
+        this.errorLogUrls.clear();
+        this.firstErrorMsg = "";
+    }
+
     protected void replayUpdateProgress(RLTaskTxnCommitAttachment attachment) {
         try {
             updateNumOfData(attachment.getTotalRows(), 
attachment.getFilteredRows(), attachment.getUnselectedRows(),
@@ -1405,8 +1413,10 @@ public abstract class RoutineLoadJob
             
routineLoadTaskInfo.handleTaskByTxnCommitAttachment(rlTaskTxnCommitAttachment);
         }
 
-        if (rlTaskTxnCommitAttachment != null && 
!Strings.isNullOrEmpty(rlTaskTxnCommitAttachment.getErrorLogUrl())) {
+        if (rlTaskTxnCommitAttachment != null
+                && 
!Strings.isNullOrEmpty(rlTaskTxnCommitAttachment.getErrorLogUrl())) {
             errorLogUrls.add(rlTaskTxnCommitAttachment.getErrorLogUrl());
+            firstErrorMsg = 
Strings.nullToEmpty(rlTaskTxnCommitAttachment.getFirstErrorMsg());
         }
 
         routineLoadTaskInfo.setTxnStatus(txnStatus);
@@ -1709,6 +1719,7 @@ public abstract class RoutineLoadJob
             row.add(userIdentity.getQualifiedUser());
             row.add(comment);
             row.add(getClusterInfo());
+            row.add(getFirstErrorMsg());
             return row;
         } finally {
             readUnlock();
@@ -1942,6 +1953,8 @@ public abstract class RoutineLoadJob
 
     @Override
     public void gsonPostProcess() throws IOException {
+        errorLogUrls = EvictingQueue.create(3);
+        firstErrorMsg = "";
         if (tableId == 0) {
             isMultiTable = true;
         }
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/ShowRoutineLoadCommand.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/ShowRoutineLoadCommand.java
index fe5c1441ebc..42737314470 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/ShowRoutineLoadCommand.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/ShowRoutineLoadCommand.java
@@ -110,6 +110,7 @@ public class ShowRoutineLoadCommand extends ShowCommand {
             .add("User")
             .add("Comment")
             .add("ComputeGroup")
+            .add("FirstErrorMsg")
             .build();
 
     private final LabelNameInfo labelNameInfo;
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/service/FrontendServiceImpl.java 
b/fe/fe-core/src/main/java/org/apache/doris/service/FrontendServiceImpl.java
index 12390e1952e..b047fbc8958 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/service/FrontendServiceImpl.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/service/FrontendServiceImpl.java
@@ -4761,6 +4761,7 @@ public class FrontendServiceImpl implements 
FrontendService.Iface {
             jobInfo.setLag(job.getLag());
             jobInfo.setReasonOfStateChanged(job.getStateReason());
             jobInfo.setErrorLogUrls(Joiner.on(", 
").join(job.getErrorLogUrls()));
+            jobInfo.setFirstErrorMsg(job.getFirstErrorMsg());
             jobInfo.setUserName(job.getUserIdentity().getQualifiedUser());
             
jobInfo.setCurrentAbortTaskNum(job.getJobStatistic().currentAbortedTaskNum);
             jobInfo.setIsAbnormalPause(job.isAbnormalPause());
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgrTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgrTest.java
index 9df6a5f3d3a..7b900960f22 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgrTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgrTest.java
@@ -39,9 +39,11 @@ import org.apache.doris.common.LabelAlreadyUsedException;
 import org.apache.doris.common.MetaNotFoundException;
 import org.apache.doris.common.QuotaExceedException;
 import org.apache.doris.common.UserException;
+import org.apache.doris.load.routineload.RLTaskTxnCommitAttachment;
 import org.apache.doris.transaction.BeginTransactionException;
 import org.apache.doris.transaction.GlobalTransactionMgrIface;
 import org.apache.doris.transaction.TransactionState;
+import org.apache.doris.transaction.TxnStateChangeCallback;
 
 import com.google.common.collect.Lists;
 import mockit.Mock;
@@ -50,6 +52,9 @@ import org.junit.Assert;
 import org.junit.Before;
 import org.junit.Test;
 import org.junit.jupiter.api.Assertions;
+import org.mockito.ArgumentCaptor;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
 
 import java.lang.reflect.InvocationTargetException;
 import java.util.concurrent.atomic.AtomicLong;
@@ -306,6 +311,53 @@ public class CloudGlobalTransactionMgrTest {
         masterTransMgr.abortTransaction(CatalogTestUtil.testDbId1, 
transactionId, "User Cancelled");
     }
 
+    @Test
+    public void testAbortRoutineLoadTransactionWithAttachment() throws 
Exception {
+        long transactionId = 123534;
+        long jobId = 1001;
+        RLTaskTxnCommitAttachment attachment = new RLTaskTxnCommitAttachment(
+                Cloud.RLTaskTxnCommitAttachmentPB.newBuilder()
+                        .setJobId(jobId)
+                        
.setTaskId(Cloud.UniqueIdPB.newBuilder().setHi(1).setLo(2))
+                        .setProgress(Cloud.RoutineLoadProgressPB.newBuilder())
+                        .setFirstErrorMsg("invalid source row")
+                        .build());
+        TxnStateChangeCallback callback = 
Mockito.mock(TxnStateChangeCallback.class);
+        Mockito.when(callback.getId()).thenReturn(jobId);
+        masterTransMgr.getCallbackFactory().addCallback(callback);
+
+        MetaServiceProxy mockProxy = Mockito.mock(MetaServiceProxy.class);
+        try (MockedStatic<MetaServiceProxy> mockedStatic = 
Mockito.mockStatic(MetaServiceProxy.class)) {
+            
mockedStatic.when(MetaServiceProxy::getInstance).thenReturn(mockProxy);
+            Mockito.doAnswer(invocation -> {
+                Cloud.AbortTxnRequest request = invocation.getArgument(0);
+                Assert.assertTrue(request.hasCommitAttachment());
+                Assert.assertEquals("invalid source row", 
request.getCommitAttachment()
+                        .getRlTaskTxnCommitAttachment().getFirstErrorMsg());
+                return AbortTxnResponse.newBuilder()
+                        .setStatus(Cloud.MetaServiceResponseStatus.newBuilder()
+                                .setCode(MetaServiceCode.OK).setMsg("OK"))
+                        .setTxnInfo(buildTxnInfo(transactionId).toBuilder()
+                                
.setStatus(Cloud.TxnStatusPB.TXN_STATUS_ABORTED)
+                                .setReason("data quality error")
+                                
.setCommitAttachment(request.getCommitAttachment()))
+                        .build();
+            }).when(mockProxy).abortTxn(Mockito.any());
+
+            masterTransMgr.abortTransaction(CatalogTestUtil.testDbId1, 
transactionId,
+                    "data quality error", attachment, Lists.newArrayList());
+
+            ArgumentCaptor<TransactionState> txnStateCaptor = 
ArgumentCaptor.forClass(TransactionState.class);
+            Mockito.verify(callback).afterAborted(txnStateCaptor.capture(), 
Mockito.eq(true),
+                    Mockito.eq("data quality error"));
+            RLTaskTxnCommitAttachment callbackAttachment =
+                    (RLTaskTxnCommitAttachment) 
txnStateCaptor.getValue().getTxnCommitAttachment();
+            Assert.assertEquals("invalid source row", 
callbackAttachment.getFirstErrorMsg());
+        } finally {
+            masterTransMgr.getCallbackFactory().removeCallback(jobId);
+        }
+    }
+
     @Test
     public void testAbortTransactionByLabel() throws UserException {
         new MockUp<MetaServiceProxy>(MetaServiceProxy.class) {
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/load/routineload/RoutineLoadJobTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/load/routineload/RoutineLoadJobTest.java
index d2246406428..cc6962560f8 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/load/routineload/RoutineLoadJobTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/load/routineload/RoutineLoadJobTest.java
@@ -21,6 +21,8 @@ import org.apache.doris.analysis.UserIdentity;
 import org.apache.doris.catalog.Database;
 import org.apache.doris.catalog.Env;
 import org.apache.doris.catalog.Table;
+import org.apache.doris.cloud.transaction.TxnUtil;
+import org.apache.doris.common.Config;
 import org.apache.doris.common.InternalErrorCode;
 import org.apache.doris.common.LoadException;
 import org.apache.doris.common.Pair;
@@ -30,6 +32,9 @@ import org.apache.doris.datasource.InternalCatalog;
 import org.apache.doris.datasource.kafka.KafkaUtil;
 import 
org.apache.doris.nereids.trees.plans.commands.info.CreateRoutineLoadInfo;
 import org.apache.doris.thrift.TKafkaRLTaskProgress;
+import org.apache.doris.thrift.TLoadSourceType;
+import org.apache.doris.thrift.TRLTaskTxnCommitAttachment;
+import org.apache.doris.thrift.TUniqueId;
 import org.apache.doris.thrift.TUniqueKeyUpdateMode;
 import org.apache.doris.transaction.TransactionException;
 import org.apache.doris.transaction.TransactionState;
@@ -51,6 +56,27 @@ import java.util.List;
 import java.util.Map;
 
 public class RoutineLoadJobTest {
+    @Test
+    public void testFirstErrorMsgInTxnCommitAttachment() {
+        String overlongFirstErrorMsg = Strings.repeat("x", 
Config.first_error_msg_max_length + 1);
+        TRLTaskTxnCommitAttachment thriftAttachment = new 
TRLTaskTxnCommitAttachment();
+        thriftAttachment.setLoadSourceType(TLoadSourceType.KAFKA);
+        thriftAttachment.setId(new TUniqueId(1, 2));
+        thriftAttachment.setJobId(3);
+        thriftAttachment.setKafkaRLTaskProgress(new 
TKafkaRLTaskProgress(Maps.newHashMap()));
+        thriftAttachment.setErrorLogUrl("http://127.0.0.1/error_log";);
+        thriftAttachment.setFirstErrorMsg(overlongFirstErrorMsg);
+
+        RLTaskTxnCommitAttachment attachment = new 
RLTaskTxnCommitAttachment(thriftAttachment);
+        Assert.assertEquals(Config.first_error_msg_max_length, 
attachment.getFirstErrorMsg().length());
+        Assert.assertTrue(attachment.getFirstErrorMsg().endsWith("..."));
+
+        RLTaskTxnCommitAttachment cloudAttachment = 
TxnUtil.rtTaskTxnCommitAttachmentFromPb(
+                TxnUtil.rlTaskTxnCommitAttachmentToPb(attachment));
+        Assert.assertEquals("http://127.0.0.1/error_log";, 
cloudAttachment.getErrorLogUrl());
+        Assert.assertEquals(attachment.getFirstErrorMsg(), 
cloudAttachment.getFirstErrorMsg());
+    }
+
     @Test
     public void testAfterAbortedReasonOffsetOutOfRange(@Mocked Env env,
                                                        @Injectable 
TransactionState transactionState,
@@ -98,6 +124,8 @@ public class RoutineLoadJobTest {
         tKafkaRLTaskProgress.partitionCmtOffset = Maps.newHashMap();
         KafkaProgress kafkaProgress = new KafkaProgress(tKafkaRLTaskProgress);
         Deencapsulation.setField(attachment, "progress", kafkaProgress);
+        Deencapsulation.setField(attachment, "errorLogUrl", 
"http://127.0.0.1/error_log";);
+        Deencapsulation.setField(attachment, "firstErrorMsg", "invalid source 
row");
 
         KafkaProgress currentProgress = new 
KafkaProgress(tKafkaRLTaskProgress);
 
@@ -134,6 +162,8 @@ public class RoutineLoadJobTest {
 
         Assert.assertEquals(RoutineLoadJob.JobState.RUNNING, 
routineLoadJob.getState());
         Assert.assertEquals(new Long(1), 
Deencapsulation.getField(jobStatistic, "abortedTaskNum"));
+        Assert.assertEquals("http://127.0.0.1/error_log";, 
routineLoadJob.getErrorLogUrls().peek());
+        Assert.assertEquals("invalid source row", 
routineLoadJob.getFirstErrorMsg());
     }
 
     @Test
@@ -184,10 +214,13 @@ public class RoutineLoadJobTest {
         Deencapsulation.setField(routineLoadJob, "pauseReason", errorReason);
         Deencapsulation.setField(routineLoadJob, "progress", kafkaProgress);
         Deencapsulation.setField(routineLoadJob, "userIdentity", userIdentity);
+        Deencapsulation.setField(routineLoadJob, "firstErrorMsg", "invalid 
source row");
 
         List<String> showInfo = routineLoadJob.getShowInfo();
         Assert.assertEquals(true, showInfo.stream().filter(entity -> 
!Strings.isNullOrEmpty(entity))
                 .anyMatch(entity -> entity.equals(errorReason.toString())));
+        Assert.assertEquals(24, showInfo.size());
+        Assert.assertEquals("invalid source row", showInfo.get(23));
     }
 
     @Test
@@ -304,11 +337,17 @@ public class RoutineLoadJobTest {
         RoutineLoadStatistic jobStatistic = 
Deencapsulation.getField(routineLoadJob, "jobStatistic");
         Deencapsulation.setField(jobStatistic, "currentErrorRows", 1);
         Deencapsulation.setField(jobStatistic, "currentTotalRows", 99);
+        Deencapsulation.setField(routineLoadJob, "otherMsg", "previous 
warning");
+        routineLoadJob.getErrorLogUrls().add("http://127.0.0.1/error_log";);
+        Deencapsulation.setField(routineLoadJob, "firstErrorMsg", "invalid 
source row");
         Deencapsulation.invoke(routineLoadJob, "updateNumOfData", 2L, 0L, 0L, 
1L, 1L, false);
 
         Assert.assertEquals(RoutineLoadJob.JobState.RUNNING, 
Deencapsulation.getField(routineLoadJob, "state"));
         Assert.assertEquals(new Long(0), 
Deencapsulation.getField(jobStatistic, "currentErrorRows"));
         Assert.assertEquals(new Long(0), 
Deencapsulation.getField(jobStatistic, "currentTotalRows"));
+        Assert.assertEquals("", Deencapsulation.getField(routineLoadJob, 
"otherMsg"));
+        Assert.assertTrue(routineLoadJob.getErrorLogUrls().isEmpty());
+        Assert.assertEquals("", routineLoadJob.getFirstErrorMsg());
 
     }
 
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/ShowRoutineLoadCommandTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/ShowRoutineLoadCommandTest.java
index 4bbe2722030..ce933be3dec 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/ShowRoutineLoadCommandTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/ShowRoutineLoadCommandTest.java
@@ -35,5 +35,7 @@ public class ShowRoutineLoadCommandTest extends 
TestWithFeService {
         LabelNameInfo labelNameInfo = new LabelNameInfo("test_db", 
"test_label");
         ShowRoutineLoadCommand command = new 
ShowRoutineLoadCommand(labelNameInfo, null, false);
         Assertions.assertDoesNotThrow(() -> command.validate(connectContext));
+        Assertions.assertEquals(24, command.getMetaData().getColumnCount());
+        Assertions.assertEquals("FirstErrorMsg", 
command.getMetaData().getColumn(23).getName());
     }
 }
diff --git a/gensrc/proto/cloud.proto b/gensrc/proto/cloud.proto
index c6a04aba85b..3cf4fccc676 100644
--- a/gensrc/proto/cloud.proto
+++ b/gensrc/proto/cloud.proto
@@ -377,6 +377,7 @@ message TxnCoordinatorPB {
 message RoutineLoadProgressPB {
     map<int32, int64> partition_to_offset = 1;
     optional RoutineLoadJobStatisticPB stat = 2;
+    // Error log URLs and first error messages are transient diagnostics and 
are not stored here.
 }
 
 message RLTaskTxnCommitAttachmentPB {
@@ -389,6 +390,7 @@ message RLTaskTxnCommitAttachmentPB {
     optional int64 task_execution_time_ms = 7;
     optional RoutineLoadProgressPB progress = 8;
     optional string error_log_url = 9;
+    optional string first_error_msg = 10;
 }
 
 message RoutineLoadJobStatisticPB {
diff --git a/gensrc/thrift/FrontendService.thrift 
b/gensrc/thrift/FrontendService.thrift
index ee2570dd576..a3f4437776d 100644
--- a/gensrc/thrift/FrontendService.thrift
+++ b/gensrc/thrift/FrontendService.thrift
@@ -643,6 +643,7 @@ struct TRLTaskTxnCommitAttachment {
     9: optional i64 loadCostMs
     10: optional TKafkaRLTaskProgress kafkaRLTaskProgress
     11: optional string errorLogUrl
+    12: optional string firstErrorMsg
 }
 
 struct TTxnCommitAttachment {
@@ -1663,6 +1664,7 @@ struct TRoutineLoadJob {
     19: optional i32 current_abort_task_num
     20: optional bool is_abnormal_pause
     21: optional string compute_group
+    22: optional string first_error_msg
 }
 
 struct TFetchRoutineLoadJobResult {
diff --git 
a/regression-test/suites/load_p0/routine_load/test_routine_load_first_error_msg.groovy
 
b/regression-test/suites/load_p0/routine_load/test_routine_load_first_error_msg.groovy
new file mode 100644
index 00000000000..a41c84d2ce0
--- /dev/null
+++ 
b/regression-test/suites/load_p0/routine_load/test_routine_load_first_error_msg.groovy
@@ -0,0 +1,127 @@
+// 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.
+
+import org.apache.doris.regression.util.RoutineLoadTestUtils
+import org.apache.kafka.clients.admin.AdminClient
+import org.apache.kafka.clients.admin.NewTopic
+import org.apache.kafka.clients.producer.ProducerConfig
+import org.junit.Assert
+
+suite("test_routine_load_first_error_msg", "p0") {
+    if (!RoutineLoadTestUtils.isKafkaTestEnabled(context)) {
+        return
+    }
+
+    def kafkaBroker = RoutineLoadTestUtils.getKafkaBroker(context)
+    def kafkaTopic = 
"test_routine_load_first_error_msg_${System.currentTimeMillis()}"
+    def jobName = "test_routine_load_first_error_msg"
+
+    def adminProps = new Properties()
+    adminProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaBroker)
+    def adminClient = AdminClient.create(adminProps)
+    try {
+        adminClient.createTopics([new NewTopic(kafkaTopic, 1, (short) 
1)]).all().get()
+    } finally {
+        adminClient.close()
+    }
+
+    try {
+        sql "DROP TABLE IF EXISTS test_routine_load_first_error_msg"
+        sql """
+            CREATE TABLE test_routine_load_first_error_msg (
+                id INT,
+                name STRING
+            )
+            PARTITION BY RANGE(id) (
+                PARTITION p0 VALUES LESS THAN ("10")
+            )
+            DISTRIBUTED BY HASH(id) BUCKETS 1
+            PROPERTIES ("replication_num" = "1")
+        """
+
+        sql """
+            CREATE ROUTINE LOAD ${jobName} ON test_routine_load_first_error_msg
+            COLUMNS TERMINATED BY "|"
+            PROPERTIES (
+                "max_error_number" = "10",
+                "max_filter_ratio" = "1.0",
+                "max_batch_interval" = "5"
+            )
+            FROM KAFKA (
+                "kafka_broker_list" = "${kafkaBroker}",
+                "kafka_topic" = "${kafkaTopic}",
+                "property.kafka_default_offsets" = "OFFSET_END"
+            )
+        """
+
+        def count = 0
+        while (true) {
+            def showResult = sql "SHOW ROUTINE LOAD FOR ${jobName}"
+            if (showResult[0][8].toString() == "RUNNING") {
+                break
+            }
+            if (count++ > 60) {
+                Assert.fail("Routine Load job did not enter RUNNING state")
+            }
+            sleep(1000)
+        }
+
+        def producer = RoutineLoadTestUtils.createKafkaProducer(kafkaBroker)
+        try {
+            RoutineLoadTestUtils.sendTestDataToKafka(
+                    producer, [kafkaTopic], ["100|bad_row", "1|valid_row"])
+            producer.flush()
+        } finally {
+            producer.close()
+        }
+
+        count = 0
+        while (true) {
+            def jobInfo = sql """
+                SELECT ERROR_LOG_URLS, FIRST_ERROR_MSG
+                FROM information_schema.routine_load_jobs
+                WHERE JOB_NAME = '${jobName}'
+            """
+            def showResult = sql "SHOW ROUTINE LOAD FOR ${jobName}"
+            def loadedRows = sql "SELECT COUNT(*) FROM 
test_routine_load_first_error_msg"
+            def informationSchemaErrorLogUrls = jobInfo[0][0].toString()
+            def firstErrorMsg = jobInfo[0][1]
+            def showErrorLogUrls = showResult[0][18].toString()
+            def showFirstErrorMsg = showResult[0][23]
+            if (loadedRows[0][0] == 1 && firstErrorMsg != null
+                    && firstErrorMsg.toString().contains("bad_row")
+                    && showFirstErrorMsg != null
+                    && showFirstErrorMsg.toString().contains("bad_row")) {
+                assertEquals(informationSchemaErrorLogUrls, showErrorLogUrls)
+                
assertFalse(informationSchemaErrorLogUrls.contains("first_error_msg:"))
+                assertFalse(showErrorLogUrls.contains("first_error_msg:"))
+                assertEquals(firstErrorMsg.toString(), 
showFirstErrorMsg.toString())
+                break
+            }
+            if (count++ > 60) {
+                Assert.fail("First error message was not reported by 
routine_load_jobs and SHOW ROUTINE LOAD: "
+                        + "loadedRows=${loadedRows[0][0]}, 
firstErrorMsg=${firstErrorMsg}, "
+                        + "showFirstErrorMsg=${showFirstErrorMsg}, "
+                        + 
"informationSchemaErrorLogUrlsEmpty=${informationSchemaErrorLogUrls.isEmpty()}, 
"
+                        + 
"showErrorLogUrlsEmpty=${showErrorLogUrls.isEmpty()}")
+            }
+            sleep(1000)
+        }
+    } finally {
+        try_sql "STOP ROUTINE LOAD FOR ${jobName}"
+    }
+}


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


Reply via email to