This is an automated email from the ASF dual-hosted git repository. yiguolei pushed a commit to branch branch-4.2 in repository https://gitbox.apache.org/repos/asf/doris.git
commit f48a453ae21812e3ac16103c22703846977d0b63 Author: meiyi <[email protected]> AuthorDate: Wed Sep 16 20:49:47 2026 +0800 branch-4.1: [fix](fe) Keep rollup jobs waiting on conflict txn abort failure (#67661) (#68055) Cherry-picked from https://github.com/apache/doris/pull/67661 --- .../java/org/apache/doris/alter/RollupJobV2.java | 10 ++++-- .../doris/transaction/GlobalTransactionMgr.java | 4 +++ .../org/apache/doris/alter/RollupJobV2Test.java | 39 ++++++++++++++++++++++ .../transaction/GlobalTransactionMgrTest.java | 21 ++++++++++++ 4 files changed, 72 insertions(+), 2 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/alter/RollupJobV2.java b/fe/fe-core/src/main/java/org/apache/doris/alter/RollupJobV2.java index 1e42d0b0382..5012eacab13 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/alter/RollupJobV2.java +++ b/fe/fe-core/src/main/java/org/apache/doris/alter/RollupJobV2.java @@ -708,8 +708,14 @@ public class RollupJobV2 extends AlterJobV2 implements GsonPostProcessable { if (Config.enable_abort_txn_by_checking_conflict_txn) { List<TransactionState> failedTxns = GlobalTransactionMgr.checkFailedTxns(unFinishedTxns); for (TransactionState txn : failedTxns) { - Env.getCurrentGlobalTransactionMgr() - .abortTransaction(txn.getDbId(), txn.getTransactionId(), "Cancel by schema change"); + try { + Env.getCurrentGlobalTransactionMgr() + .abortTransaction(txn.getDbId(), txn.getTransactionId(), "Cancel by schema change"); + } catch (UserException e) { + LOG.warn("failed to abort previous load txn {}, wait next round. rollup job: {}", + txn.getTransactionId(), jobId, e); + return false; + } } } return unFinishedTxns.isEmpty(); diff --git a/fe/fe-core/src/main/java/org/apache/doris/transaction/GlobalTransactionMgr.java b/fe/fe-core/src/main/java/org/apache/doris/transaction/GlobalTransactionMgr.java index 2ce7717912f..bd1ddf2ba6a 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/transaction/GlobalTransactionMgr.java +++ b/fe/fe-core/src/main/java/org/apache/doris/transaction/GlobalTransactionMgr.java @@ -505,6 +505,10 @@ public class GlobalTransactionMgr implements GlobalTransactionMgrIface { public static List<TransactionState> checkFailedTxns(List<TransactionState> conflictTxns) { List<TransactionState> failedTxns = new ArrayList<>(); for (TransactionState txn : conflictTxns) { + TransactionStatus status = txn.getTransactionStatus(); + if (status == TransactionStatus.COMMITTED || status.isFinalStatus()) { + continue; + } if (checkFailedTxnsByCoordinator(txn)) { failedTxns.add(txn); } diff --git a/fe/fe-core/src/test/java/org/apache/doris/alter/RollupJobV2Test.java b/fe/fe-core/src/test/java/org/apache/doris/alter/RollupJobV2Test.java index 38ae8944fc8..ff2772e9482 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/alter/RollupJobV2Test.java +++ b/fe/fe-core/src/test/java/org/apache/doris/alter/RollupJobV2Test.java @@ -52,7 +52,13 @@ import org.apache.doris.task.AgentTaskQueue; import org.apache.doris.thrift.TStorageFormat; import org.apache.doris.thrift.TTaskType; import org.apache.doris.transaction.FakeTransactionIDGenerator; +import org.apache.doris.transaction.GlobalTransactionMgr; import org.apache.doris.transaction.GlobalTransactionMgrIface; +import org.apache.doris.transaction.TransactionState; +import org.apache.doris.transaction.TransactionState.LoadJobSourceType; +import org.apache.doris.transaction.TransactionState.TxnCoordinator; +import org.apache.doris.transaction.TransactionState.TxnSourceType; +import org.apache.doris.transaction.TransactionStatus; import com.google.common.collect.Lists; import mockit.Mock; @@ -61,6 +67,8 @@ import org.junit.After; import org.junit.Assert; import org.junit.Before; import org.junit.Test; +import org.mockito.MockedStatic; +import org.mockito.Mockito; import java.io.DataInputStream; import java.io.DataOutputStream; @@ -130,6 +138,37 @@ public class RollupJobV2Test { file.delete(); } + @Test + public void testCommitWhileAbortingPreviousLoad() throws Exception { + long txnId = masterTransMgr.beginTransaction(CatalogTestUtil.testDbId1, + Lists.newArrayList(CatalogTestUtil.testTableId1), "commit_during_rollup_abort", + new TxnCoordinator(TxnSourceType.FE, 0, "missing", 0), LoadJobSourceType.FRONTEND, 60); + TransactionState txn = masterTransMgr.getTransactionState(CatalogTestUtil.testDbId1, txnId); + Database db = masterEnv.getInternalCatalog().getDbOrDdlException(CatalogTestUtil.testDbId1); + OlapTable table = (OlapTable) db.getTableOrDdlException(CatalogTestUtil.testTableId1); + MaterializedViewHandler handler = masterEnv.getMaterializedViewHandler(); + handler.process(Lists.newArrayList(clause), db, table); + RollupJobV2 job = (RollupJobV2) handler.getAlterJobsV2().values().iterator().next(); + job.jobState = JobState.WAITING_TXN; + job.watershedTxnId = txnId + 1; + + try (MockedStatic<GlobalTransactionMgr> mocked = Mockito.mockStatic( + GlobalTransactionMgr.class, Mockito.CALLS_REAL_METHODS)) { + mocked.when(() -> GlobalTransactionMgr.checkFailedTxns(Mockito.anyList())).thenAnswer(invocation -> { + List<TransactionState> failed = (List<TransactionState>) invocation.callRealMethod(); + Assert.assertEquals(Lists.newArrayList(txn), failed); + txn.setTransactionStatus(TransactionStatus.COMMITTED); + return failed; + }); + job.runWaitingTxnJob(); + Assert.assertEquals(JobState.WAITING_TXN, job.getJobState()); + Assert.assertEquals(TransactionStatus.COMMITTED, txn.getTransactionStatus()); + } + Assert.assertFalse(job.checkFailedPreviousLoadAndAbort()); + txn.setTransactionStatus(TransactionStatus.VISIBLE); + Assert.assertTrue(job.checkFailedPreviousLoadAndAbort()); + } + @Test public void testRunRollupJobConcurrentLimit() throws UserException { fakeEnv = new FakeEnv(); diff --git a/fe/fe-core/src/test/java/org/apache/doris/transaction/GlobalTransactionMgrTest.java b/fe/fe-core/src/test/java/org/apache/doris/transaction/GlobalTransactionMgrTest.java index 10e446348af..2f623c0022a 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/transaction/GlobalTransactionMgrTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/transaction/GlobalTransactionMgrTest.java @@ -116,6 +116,27 @@ public class GlobalTransactionMgrTest { slaveTransMgr.setEditLog(slaveEnv.getEditLog()); } + @Test + public void testCheckFailedTxnsWithOfflineCoordinator() { + FakeEnv.setEnv(masterEnv); + for (TxnSourceType source : Lists.newArrayList(TxnSourceType.FE, TxnSourceType.BE)) { + TransactionState txn = new TransactionState(CatalogTestUtil.testDbId1, + Lists.newArrayList(CatalogTestUtil.testTableId1), 1, "offline_coordinator", null, + LoadJobSourceType.FRONTEND, new TxnCoordinator(source, Long.MAX_VALUE, "missing", 0), -1, 60000); + for (TransactionStatus status : Lists.newArrayList( + TransactionStatus.PREPARE, TransactionStatus.PRECOMMITTED)) { + txn.setTransactionStatus(status); + Assert.assertEquals(Lists.newArrayList(txn), + GlobalTransactionMgr.checkFailedTxns(Lists.newArrayList(txn))); + } + for (TransactionStatus status : Lists.newArrayList(TransactionStatus.COMMITTED, + TransactionStatus.VISIBLE, TransactionStatus.ABORTED)) { + txn.setTransactionStatus(status); + Assert.assertTrue(GlobalTransactionMgr.checkFailedTxns(Lists.newArrayList(txn)).isEmpty()); + } + } + } + @Test public void testBeginTransaction() throws LabelAlreadyUsedException, AnalysisException, BeginTransactionException, DuplicatedRequestException, QuotaExceedException, MetaNotFoundException { --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
