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 25564487d3f branch-4.1: [fix](job) Recover broker load from visible 
transaction #66987 (#67075)
25564487d3f is described below

commit 25564487d3f5d10f7831ead9b1d9cfd18e0b43c5
Author: github-actions[bot] 
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Tue Aug 25 10:40:07 2026 +0800

    branch-4.1: [fix](job) Recover broker load from visible transaction #66987 
(#67075)
    
    Cherry-picked from #66987
    
    Co-authored-by: hui lai <[email protected]>
---
 .../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 4fb3ac88911..feeaddb0edc 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
@@ -150,27 +150,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 767c360c880..135a7485850 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
@@ -23,6 +23,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.jmockit.Deencapsulation;
@@ -427,10 +428,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);
 
@@ -440,6 +441,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.
@@ -456,8 +497,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