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]