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

liaoxin01 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 5f72a92e621 [fix](job) Recover broker load from visible transaction 
(#66987)
5f72a92e621 is described below

commit 5f72a92e62129f2cac083e24d5542175dc650bf7
Author: hui lai <[email protected]>
AuthorDate: Mon Aug 24 19:56:43 2026 +0800

    [fix](job) Recover broker load from visible transaction (#66987)
    
    ### What problem does this PR solve?
    
    In cloud mode, a Broker Load transaction can become VISIBLE in Meta
    Service before FE persists the final Load Job record. If FE restarts in
    that window, the replayed PENDING job begins the same label again and is
    cancelled after conflicting with its own visible transaction.
    
    This change looks up the transaction by label, verifies that its
    callback ID belongs to the current Load Job, reuses PREPARE
    transactions, and finishes the job immediately when the transaction is
    already VISIBLE. For Cloud Broker Load, the final operation is persisted
    to EditLog.
    
    ### Release note
    
    Fix Broker Load jobs being reported as cancelled after FE restart even
    though their transactions were already visible.
---
 .../apache/doris/load/loadv2/BrokerLoadJob.java    | 26 +++++++-----
 .../doris/load/loadv2/BrokerLoadJobTest.java       | 48 ++++++++++++++++++++--
 2 files changed, 60 insertions(+), 14 deletions(-)

diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/load/loadv2/BrokerLoadJob.java 
b/fe/fe-core/src/main/java/org/apache/doris/load/loadv2/BrokerLoadJob.java
index 5add6d9a9b4..3c43e26f971 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/load/loadv2/BrokerLoadJob.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/load/loadv2/BrokerLoadJob.java
@@ -154,27 +154,33 @@ public class BrokerLoadJob extends BulkLoadJob {
             // before transactionId was assigned (e.g. edit log write 
failure), and the pending
             // task retry would otherwise burn all retries on this exception 
and cancel the job
             // with a misleading "Label has already been used".
-            Long ownTxnId = findSelfPreparedTxnByLabel();
-            if (ownTxnId == null) {
+            TransactionState ownTxn = findSelfTxnByLabel();
+            if (ownTxn == null) {
                 throw e;
             }
-            LOG.info("broker load job {} adopts its own prepared txn {} for 
label {} on pending task retry",
-                    id, ownTxnId, label);
-            transactionId = ownTxnId;
+            transactionId = ownTxn.getTransactionId();
+            if (ownTxn.getTransactionStatus() == TransactionStatus.VISIBLE) {
+                LOG.info("broker load job {} recovers its own visible txn {} 
for label {}",
+                        id, transactionId, label);
+                afterVisible(ownTxn, true);
+            } else {
+                LOG.info("broker load job {} adopts its own prepared txn {} 
for label {} on pending task retry",
+                        id, transactionId, label);
+            }
         }
     }
 
-    private Long findSelfPreparedTxnByLabel() {
+    private TransactionState findSelfTxnByLabel() {
         try {
-            Long txnId = 
Env.getCurrentGlobalTransactionMgr().getTransactionIdByLabel(dbId, label,
-                    Lists.newArrayList(TransactionStatus.PREPARE));
+            Long txnId = 
Env.getCurrentGlobalTransactionMgr().getTransactionId(dbId, label);
             if (txnId == null) {
                 return null;
             }
             TransactionState existingTxn = 
Env.getCurrentGlobalTransactionMgr().getTransactionState(dbId, txnId);
             if (existingTxn != null && existingTxn.getCallbackId() == id
-                    && existingTxn.getTransactionStatus() == 
TransactionStatus.PREPARE) {
-                return txnId;
+                    && (existingTxn.getTransactionStatus() == 
TransactionStatus.PREPARE
+                        || existingTxn.getTransactionStatus() == 
TransactionStatus.VISIBLE)) {
+                return existingTxn;
             }
         } catch (Exception lookupException) {
             LOG.warn("broker load job {} failed to look up txn by label {}", 
id, label, lookupException);
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/load/loadv2/BrokerLoadJobTest.java 
b/fe/fe-core/src/test/java/org/apache/doris/load/loadv2/BrokerLoadJobTest.java
index efe785f2998..21dcba9e988 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/load/loadv2/BrokerLoadJobTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/load/loadv2/BrokerLoadJobTest.java
@@ -24,6 +24,7 @@ import org.apache.doris.catalog.Env;
 import org.apache.doris.catalog.OlapTable;
 import org.apache.doris.catalog.Table;
 import org.apache.doris.catalog.TableProperty;
+import org.apache.doris.cloud.load.CloudBrokerLoadJob;
 import org.apache.doris.common.LabelAlreadyUsedException;
 import org.apache.doris.common.MetaNotFoundException;
 import org.apache.doris.common.Status;
@@ -499,10 +500,10 @@ public class BrokerLoadJobTest {
                     Mockito.anyString(), Mockito.any(), Mockito.any(), 
Mockito.any(),
                     Mockito.anyLong(), Mockito.anyLong()))
                     .thenThrow(new 
LabelAlreadyUsedException("label_self_conflict"));
-            
Mockito.when(transactionMgr.getTransactionIdByLabel(Mockito.anyLong(), 
Mockito.anyString(),
-                    Mockito.anyList())).thenReturn(777L);
+            Mockito.when(transactionMgr.getTransactionId(Mockito.anyLong(), 
Mockito.anyString())).thenReturn(777L);
             Mockito.when(transactionMgr.getTransactionState(Mockito.anyLong(), 
Mockito.eq(777L)))
                     .thenReturn(preparedTxn);
+            Mockito.when(preparedTxn.getTransactionId()).thenReturn(777L);
             Mockito.when(preparedTxn.getCallbackId()).thenReturn(1001L);
             
Mockito.when(preparedTxn.getTransactionStatus()).thenReturn(TransactionStatus.PREPARE);
 
@@ -512,6 +513,46 @@ public class BrokerLoadJobTest {
         Assert.assertEquals(777L, (long) 
Deencapsulation.getField(brokerLoadJob, "transactionId"));
     }
 
+    @Test
+    public void testBeginTxnFinishesOwnVisibleTxn() throws Exception {
+        // The transaction may become visible in meta service before cloud FE 
persists the final
+        // load-job state. After FE restart the replayed PENDING job must 
recover that successful
+        // transaction instead of retrying beginTxn until it is cancelled by 
its own label.
+        GlobalTransactionMgrIface transactionMgr = 
Mockito.mock(GlobalTransactionMgrIface.class);
+        TransactionState visibleTxn = Mockito.mock(TransactionState.class);
+        TxnStateCallbackFactory callbackFactory = 
Mockito.mock(TxnStateCallbackFactory.class);
+        Env env = Mockito.mock(Env.class);
+        EditLog editLog = Mockito.mock(EditLog.class);
+        CloudBrokerLoadJob brokerLoadJob = new CloudBrokerLoadJob();
+        Deencapsulation.setField(brokerLoadJob, "id", 1001L);
+        Deencapsulation.setField(brokerLoadJob, "dbId", 1L);
+        Deencapsulation.setField(brokerLoadJob, "label", 
"label_visible_after_restart");
+
+        try (MockedStatic<Env> envMockedStatic = 
Mockito.mockStatic(Env.class)) {
+            
envMockedStatic.when(Env::getCurrentGlobalTransactionMgr).thenReturn(transactionMgr);
+            envMockedStatic.when(Env::getCurrentEnv).thenReturn(env);
+            Mockito.when(env.getEditLog()).thenReturn(editLog);
+            
Mockito.when(transactionMgr.getCallbackFactory()).thenReturn(callbackFactory);
+            Mockito.when(transactionMgr.beginTransaction(Mockito.anyLong(), 
Mockito.anyList(),
+                    Mockito.anyString(), Mockito.any(), Mockito.any(), 
Mockito.any(),
+                    Mockito.anyLong(), Mockito.anyLong()))
+                    .thenThrow(new 
LabelAlreadyUsedException("label_visible_after_restart"));
+            Mockito.when(transactionMgr.getTransactionId(Mockito.anyLong(), 
Mockito.anyString())).thenReturn(888L);
+            Mockito.when(transactionMgr.getTransactionState(Mockito.anyLong(), 
Mockito.eq(888L)))
+                    .thenReturn(visibleTxn);
+            Mockito.when(visibleTxn.getTransactionId()).thenReturn(888L);
+            Mockito.when(visibleTxn.getCallbackId()).thenReturn(1001L);
+            
Mockito.when(visibleTxn.getTransactionStatus()).thenReturn(TransactionStatus.VISIBLE);
+
+            brokerLoadJob.beginTxn();
+        }
+
+        Assert.assertEquals(888L, (long) 
Deencapsulation.getField(brokerLoadJob, "transactionId"));
+        Assert.assertEquals(JobState.FINISHED, brokerLoadJob.getState());
+        Mockito.verify(callbackFactory).removeCallback(1001L);
+        
Mockito.verify(editLog).logEndLoadJob(Mockito.any(LoadJobFinalOperation.class));
+    }
+
     @Test
     public void testBeginTxnRethrowsForeignLabelConflict() throws Exception {
         // The label belongs to some other job's txn: the original exception 
must propagate.
@@ -528,8 +569,7 @@ public class BrokerLoadJobTest {
                     Mockito.anyString(), Mockito.any(), Mockito.any(), 
Mockito.any(),
                     Mockito.anyLong(), Mockito.anyLong()))
                     .thenThrow(new LabelAlreadyUsedException("label_foreign"));
-            
Mockito.when(transactionMgr.getTransactionIdByLabel(Mockito.anyLong(), 
Mockito.anyString(),
-                    Mockito.anyList())).thenReturn(888L);
+            Mockito.when(transactionMgr.getTransactionId(Mockito.anyLong(), 
Mockito.anyString())).thenReturn(888L);
             Mockito.when(transactionMgr.getTransactionState(Mockito.anyLong(), 
Mockito.eq(888L)))
                     .thenReturn(foreignTxn);
             Mockito.when(foreignTxn.getCallbackId()).thenReturn(9999L);


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

Reply via email to