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

gavinchou 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 b1d21a66ebd [fix](job) Prevent stale cloud broker load retries (#67250)
b1d21a66ebd is described below

commit b1d21a66ebd2af37259a3c1111b1ba5dfe3d19ff
Author: Refrain <[email protected]>
AuthorDate: Wed Sep 2 22:56:12 2026 +0800

    [fix](job) Prevent stale cloud broker load retries (#67250)
    
    Related PR: #66987
    
    Problem Summary:
    
    In cloud mode, a Broker Load attempt can have many loading tasks queued
    in the FE executor. After the first task failure moves the job to RETRY,
    a queued task can currently move it back to LOADING. This allows several
    concurrent failures to enter the job-level retry path and wait on the
    shared idToTasks map.
    
    An older retry handler can consequently follow tasks from a newer
    attempt. If that newer attempt finishes successfully, the stale handler
    can overwrite FINISHED with PENDING, schedule another pending task, and
    reuse the already VISIBLE transaction. New rowsets are then rejected by
    Meta Service because the transaction is no longer PREPARED.
    
    This change closes the current attempt once the job enters RETRY and
    verifies that the job is still RETRY before scheduling the next PENDING
    attempt. It does not change transaction protocols, Meta Service
    behavior, or the non-cloud retry implementation.
    
    ### Release note
    
    Fix Cloud Broker Load jobs being retried after a later attempt already
    finished.
---
 .../doris/cloud/load/CloudBrokerLoadJob.java       |  8 +++++
 .../java/org/apache/doris/load/loadv2/LoadJob.java |  5 ++++
 .../doris/cloud/load/CloudBrokerLoadJobTest.java   | 35 ++++++++++++++++++++++
 .../org/apache/doris/load/loadv2/LoadJobTest.java  |  9 ++++++
 4 files changed, 57 insertions(+)

diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/cloud/load/CloudBrokerLoadJob.java 
b/fe/fe-core/src/main/java/org/apache/doris/cloud/load/CloudBrokerLoadJob.java
index 1d58b6eba7f..1e1b4ff2ade 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/cloud/load/CloudBrokerLoadJob.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/cloud/load/CloudBrokerLoadJob.java
@@ -277,6 +277,14 @@ public class CloudBrokerLoadJob extends BrokerLoadJob {
 
         try {
             writeLock();
+            // The transaction callback or another terminal path may finish 
the job while tasks drain.
+            if (state != JobState.RETRY) {
+                LOG.info(new LogBuilder(LogKey.LOAD_JOB, id)
+                        .add("state", state)
+                        .add("msg", "skip retry because the load job is no 
longer retrying")
+                        .build());
+                return;
+            }
             this.state = JobState.PENDING;
             this.idToTasks.clear();
             this.failMsg = null;
diff --git a/fe/fe-core/src/main/java/org/apache/doris/load/loadv2/LoadJob.java 
b/fe/fe-core/src/main/java/org/apache/doris/load/loadv2/LoadJob.java
index 1ca2eb05562..1a1955be5e0 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/load/loadv2/LoadJob.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/load/loadv2/LoadJob.java
@@ -430,6 +430,11 @@ public abstract class LoadJob extends 
AbstractTxnStateChangeCallback
                     id, this.state, jobState);
             return false;
         }
+        // RETRY closes the current attempt. The next attempt must enter 
PENDING before LOADING.
+        if (this.state == JobState.RETRY && jobState == JobState.LOADING) {
+            LOG.info("the load job {} is retrying, should not update state to 
{}", id, jobState);
+            return false;
+        }
         switch (jobState) {
             case UNKNOWN:
                 executeUnknown();
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/cloud/load/CloudBrokerLoadJobTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/cloud/load/CloudBrokerLoadJobTest.java
index 4ef64019638..65f53757eba 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/cloud/load/CloudBrokerLoadJobTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/cloud/load/CloudBrokerLoadJobTest.java
@@ -23,6 +23,7 @@ import org.apache.doris.common.jmockit.Deencapsulation;
 import org.apache.doris.load.BrokerFileGroupAggInfo;
 import org.apache.doris.load.FailMsg;
 import org.apache.doris.load.loadv2.JobState;
+import org.apache.doris.load.loadv2.LoadTask;
 import org.apache.doris.transaction.GlobalTransactionMgrIface;
 import org.apache.doris.transaction.TxnStateCallbackFactory;
 
@@ -32,6 +33,8 @@ import org.junit.Test;
 import org.mockito.MockedStatic;
 import org.mockito.Mockito;
 
+import java.util.Map;
+
 public class CloudBrokerLoadJobTest {
 
     @Test
@@ -86,4 +89,36 @@ public class CloudBrokerLoadJobTest {
         Assert.assertEquals(0L, job.getTransactionId());
         Assert.assertEquals(JobState.RETRY, job.getState());
     }
+
+    @Test
+    public void testStaleRetryDoesNotReviveFinishedJob() throws Exception {
+        GlobalTransactionMgrIface transactionMgr = 
Mockito.mock(GlobalTransactionMgrIface.class);
+        TxnStateCallbackFactory callbackFactory = 
Mockito.mock(TxnStateCallbackFactory.class);
+        CloudBrokerLoadJob job = new CloudBrokerLoadJob();
+        LoadTask failedTask = Mockito.mock(LoadTask.class);
+        LoadTask otherTask = Mockito.mock(LoadTask.class);
+        Deencapsulation.setField(job, "id", 1003L);
+        Deencapsulation.setField(job, "dbId", 2003L);
+        Deencapsulation.setField(job, "label", 
"cloud_broker_load_finished_during_retry");
+        Deencapsulation.setField(job, "transactionId", 3003L);
+        Deencapsulation.setField(job, "cloudClusterId", "cluster_id");
+        Map<Long, LoadTask> idToTasks = Deencapsulation.getField(job, 
"idToTasks");
+        idToTasks.put(4003L, failedTask);
+        idToTasks.put(4004L, otherTask);
+        Mockito.when(otherTask.isDone()).thenAnswer(invocation -> {
+            Deencapsulation.setField(job, "state", JobState.FINISHED);
+            return true;
+        });
+
+        try (MockedStatic<Env> envMockedStatic = 
Mockito.mockStatic(Env.class)) {
+            
envMockedStatic.when(Env::getCurrentGlobalTransactionMgr).thenReturn(transactionMgr);
+            
Mockito.when(transactionMgr.getCallbackFactory()).thenReturn(callbackFactory);
+
+            job.onTaskFailed(4003L, new 
FailMsg(FailMsg.CancelType.ETL_RUN_FAIL, "rpc failed"));
+        }
+
+        Assert.assertEquals(JobState.FINISHED, job.getState());
+        Mockito.verify(callbackFactory).removeCallback(1003L);
+        Mockito.verify(callbackFactory, Mockito.never()).addCallback(job);
+    }
 }
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/load/loadv2/LoadJobTest.java 
b/fe/fe-core/src/test/java/org/apache/doris/load/loadv2/LoadJobTest.java
index 886e166fffe..e1788fa89e4 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/load/loadv2/LoadJobTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/load/loadv2/LoadJobTest.java
@@ -155,6 +155,15 @@ public class LoadJobTest {
         Assert.assertNotEquals(-1, (long) Deencapsulation.getField(loadJob, 
"loadStartTimestamp"));
     }
 
+    @Test
+    public void testRetryStateCannotReturnToLoading() {
+        LoadJob loadJob = new BrokerLoadJob();
+        Deencapsulation.setField(loadJob, "state", JobState.RETRY);
+
+        Assert.assertFalse(loadJob.updateState(JobState.LOADING));
+        Assert.assertEquals(JobState.RETRY, loadJob.getState());
+    }
+
     @Test
     public void testUpdateStateToFinished() {
         try (MockedStatic<Env> envMockedStatic = 
Mockito.mockStatic(Env.class)) {


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

Reply via email to