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 5fcda491a19 branch-4.1: [fix](fe) Preserve active export jobs and
label mappings during cleanup (#67539) (#67833)
5fcda491a19 is described below
commit 5fcda491a190cf228e81abeabb9e23d467257054
Author: Yixuan Wang <[email protected]>
AuthorDate: Sat Sep 12 00:48:16 2026 +0800
branch-4.1: [fix](fe) Preserve active export jobs and label mappings during
cleanup (#67539) (#67833)
pick: https://github.com/apache/doris/pull/67539
---
.../main/java/org/apache/doris/load/ExportMgr.java | 25 ++--
.../apache/doris/load/loadv2/ExportMgrTest.java | 142 +++++++++++++++++++++
2 files changed, 159 insertions(+), 8 deletions(-)
diff --git a/fe/fe-core/src/main/java/org/apache/doris/load/ExportMgr.java
b/fe/fe-core/src/main/java/org/apache/doris/load/ExportMgr.java
index c2b1be2a230..6c9e3570235 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/load/ExportMgr.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/load/ExportMgr.java
@@ -462,7 +462,7 @@ public class ExportMgr {
iter.remove();
Map<String, Long> labelJobs =
dbTolabelToExportJobId.get(job.getDbId());
if (labelJobs != null) {
- labelJobs.remove(job.getLabel());
+ labelJobs.remove(job.getLabel(), job.getId());
if (labelJobs.isEmpty()) {
dbTolabelToExportJobId.remove(job.getDbId());
}
@@ -473,15 +473,19 @@ public class ExportMgr {
if (exportIdToJob.size() > Config.max_export_history_job_num) {
List<Map.Entry<Long, ExportJob>> jobList = new
ArrayList<>(exportIdToJob.entrySet());
jobList.sort(Comparator.comparingLong(entry ->
entry.getValue().getCreateTimeMs()));
- while (exportIdToJob.size() >
Config.max_export_history_job_num) {
- // Remove the oldest job
- Map.Entry<Long, ExportJob> oldestEntry = jobList.remove(0);
+ Iterator<Map.Entry<Long, ExportJob>> jobIterator =
jobList.iterator();
+ while (exportIdToJob.size() >
Config.max_export_history_job_num && jobIterator.hasNext()) {
+ Map.Entry<Long, ExportJob> oldestEntry =
jobIterator.next();
+ ExportJob job = oldestEntry.getValue();
+ if (job.getState() != ExportJobState.CANCELLED &&
job.getState() != ExportJobState.FINISHED) {
+ continue;
+ }
exportIdToJob.remove(oldestEntry.getKey());
- Map<String, Long> labelJobs =
dbTolabelToExportJobId.get(oldestEntry.getValue().getDbId());
+ Map<String, Long> labelJobs =
dbTolabelToExportJobId.get(job.getDbId());
if (labelJobs != null) {
- labelJobs.remove(oldestEntry.getValue().getLabel());
+ labelJobs.remove(job.getLabel(), job.getId());
if (labelJobs.isEmpty()) {
-
dbTolabelToExportJobId.remove(oldestEntry.getValue().getDbId());
+ dbTolabelToExportJobId.remove(job.getDbId());
}
}
}
@@ -503,8 +507,13 @@ public class ExportMgr {
public void replayUpdateJobState(ExportJobStateTransfer stateTransfer) {
writeLock();
try {
- LOG.info("replay update export job: {}, {}",
stateTransfer.getJobId(), stateTransfer.getState());
ExportJob job = exportIdToJob.get(stateTransfer.getJobId());
+ if (job == null) {
+ LOG.warn("ignore replay update for missing export job: {}, {}",
+ stateTransfer.getJobId(), stateTransfer.getState());
+ return;
+ }
+ LOG.info("replay update export job: {}, {}",
stateTransfer.getJobId(), stateTransfer.getState());
job.replayExportJobState(stateTransfer.getState());
job.setStartTimeMs(stateTransfer.getStartTimeMs());
job.setFinishTimeMs(stateTransfer.getFinishTimeMs());
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/load/loadv2/ExportMgrTest.java
b/fe/fe-core/src/test/java/org/apache/doris/load/loadv2/ExportMgrTest.java
index 38549537fbc..8d91ff090f5 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/load/loadv2/ExportMgrTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/load/loadv2/ExportMgrTest.java
@@ -19,19 +19,26 @@ package org.apache.doris.load.loadv2;
import org.apache.doris.analysis.BrokerDesc;
import org.apache.doris.analysis.TableName;
+import org.apache.doris.catalog.Env;
import org.apache.doris.common.Config;
+import org.apache.doris.common.LabelAlreadyUsedException;
import org.apache.doris.common.jmockit.Deencapsulation;
import org.apache.doris.datasource.InternalCatalog;
import org.apache.doris.load.ExportJob;
import org.apache.doris.load.ExportJobState;
+import org.apache.doris.load.ExportJobStateTransfer;
import org.apache.doris.load.ExportMgr;
import org.apache.doris.mysql.privilege.AccessControllerManager;
import org.apache.doris.mysql.privilege.MockedAuth;
+import org.apache.doris.persist.EditLog;
import mockit.Mocked;
import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;
+import org.junit.jupiter.api.Assertions;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
import java.util.Comparator;
import java.util.HashMap;
@@ -122,6 +129,141 @@ public class ExportMgrTest {
}
}
+ @Test
+ public void testRemoveOldExportJobsKeepsRunningJobs() {
+ ExportMgr isolatedExportMgr = new ExportMgr();
+ int originalMaxHistoryJobNum = Config.max_export_history_job_num;
+ Config.max_export_history_job_num = 2;
+ try {
+ long currentTime = System.currentTimeMillis();
+ ExportJob pendingJob = makeExportJob(1001, "pending");
+ Deencapsulation.setField(pendingJob, "createTimeMs", currentTime -
4000);
+ Deencapsulation.setField(pendingJob, "state",
ExportJobState.PENDING);
+ isolatedExportMgr.unprotectAddJob(pendingJob);
+
+ ExportJob exportingJob = makeExportJob(1002, "exporting");
+ Deencapsulation.setField(exportingJob, "createTimeMs", currentTime
- 3000);
+ Deencapsulation.setField(exportingJob, "state",
ExportJobState.EXPORTING);
+ isolatedExportMgr.unprotectAddJob(exportingJob);
+
+ ExportJob finishedJob = makeExportJob(1003, "finished");
+ Deencapsulation.setField(finishedJob, "createTimeMs", currentTime
- 2000);
+ Deencapsulation.setField(finishedJob, "state",
ExportJobState.FINISHED);
+ isolatedExportMgr.unprotectAddJob(finishedJob);
+
+ ExportJob cancelledJob = makeExportJob(1004, "cancelled");
+ Deencapsulation.setField(cancelledJob, "createTimeMs", currentTime
- 1000);
+ Deencapsulation.setField(cancelledJob, "state",
ExportJobState.CANCELLED);
+ isolatedExportMgr.unprotectAddJob(cancelledJob);
+
+ isolatedExportMgr.removeOldExportJobs();
+
+ Assertions.assertEquals(2, isolatedExportMgr.getJobs().size());
+
Assertions.assertNotNull(isolatedExportMgr.getJob(pendingJob.getId()));
+
Assertions.assertNotNull(isolatedExportMgr.getJob(exportingJob.getId()));
+
Assertions.assertNull(isolatedExportMgr.getJob(finishedJob.getId()));
+
Assertions.assertNull(isolatedExportMgr.getJob(cancelledJob.getId()));
+ } finally {
+ Config.max_export_history_job_num = originalMaxHistoryJobNum;
+ }
+ }
+
+ @Test
+ public void testRemoveOldExportJobsKeepsRunningJobsWhenOverLimit() {
+ ExportMgr isolatedExportMgr = new ExportMgr();
+ int originalMaxHistoryJobNum = Config.max_export_history_job_num;
+ Config.max_export_history_job_num = 1;
+ try {
+ long currentTime = System.currentTimeMillis();
+ ExportJob pendingJob = makeExportJob(2001, "pending-over-limit");
+ Deencapsulation.setField(pendingJob, "createTimeMs", currentTime -
3000);
+ Deencapsulation.setField(pendingJob, "state",
ExportJobState.PENDING);
+ isolatedExportMgr.unprotectAddJob(pendingJob);
+
+ ExportJob exportingJob = makeExportJob(2002,
"exporting-over-limit");
+ Deencapsulation.setField(exportingJob, "createTimeMs", currentTime
- 2000);
+ Deencapsulation.setField(exportingJob, "state",
ExportJobState.EXPORTING);
+ isolatedExportMgr.unprotectAddJob(exportingJob);
+
+ ExportJob inQueueJob = makeExportJob(2003, "in-queue-over-limit");
+ Deencapsulation.setField(inQueueJob, "createTimeMs", currentTime -
1000);
+ Deencapsulation.setField(inQueueJob, "state",
ExportJobState.IN_QUEUE);
+ isolatedExportMgr.unprotectAddJob(inQueueJob);
+
+ isolatedExportMgr.removeOldExportJobs();
+
+ Assertions.assertEquals(3, isolatedExportMgr.getJobs().size());
+ Assertions.assertEquals(ExportJobState.PENDING,
isolatedExportMgr.getJob(pendingJob.getId()).getState());
+ Assertions.assertEquals(ExportJobState.EXPORTING,
isolatedExportMgr.getJob(exportingJob.getId()).getState());
+ Assertions.assertEquals(ExportJobState.IN_QUEUE,
isolatedExportMgr.getJob(inQueueJob.getId()).getState());
+ } finally {
+ Config.max_export_history_job_num = originalMaxHistoryJobNum;
+ }
+ }
+
+ @Test
+ public void testRemoveOldExportJobsKeepsSameLabelJobMapping() throws
Exception {
+ Env env = Mockito.mock(Env.class);
+ EditLog editLog = Mockito.mock(EditLog.class);
+ Mockito.when(env.getEditLog()).thenReturn(editLog);
+ int originalMaxHistoryJobNum = Config.max_export_history_job_num;
+ try (MockedStatic<Env> envStatic = Mockito.mockStatic(Env.class)) {
+ envStatic.when(Env::getCurrentEnv).thenReturn(env);
+
+ ExportMgr expiredJobMgr = new ExportMgr();
+ long currentTime = System.currentTimeMillis();
+ ExportJob expiredCancelledJob = makeExportJob(3001, "same-label");
+ Deencapsulation.setField(expiredCancelledJob, "createTimeMs",
+ currentTime - (Config.history_job_keep_max_second + 1) *
1000L);
+ Deencapsulation.setField(expiredCancelledJob, "state",
ExportJobState.CANCELLED);
+ expiredJobMgr.unprotectAddJob(expiredCancelledJob);
+
+ ExportJob runningJob = makeExportJob(3002, "same-label");
+ Deencapsulation.setField(runningJob, "state",
ExportJobState.EXPORTING);
+ expiredJobMgr.unprotectAddJob(runningJob);
+ expiredJobMgr.removeOldExportJobs();
+
+ assertSameLabelIsStillUsed(expiredJobMgr);
+
+ ExportMgr overLimitJobMgr = new ExportMgr();
+ Config.max_export_history_job_num = 1;
+ ExportJob overLimitCancelledJob = makeExportJob(4001,
"same-label");
+ Deencapsulation.setField(overLimitCancelledJob, "createTimeMs",
currentTime - 2000);
+ Deencapsulation.setField(overLimitCancelledJob, "state",
ExportJobState.CANCELLED);
+ overLimitJobMgr.unprotectAddJob(overLimitCancelledJob);
+
+ ExportJob overLimitRunningJob = makeExportJob(4002, "same-label");
+ Deencapsulation.setField(overLimitRunningJob, "createTimeMs",
currentTime - 1000);
+ Deencapsulation.setField(overLimitRunningJob, "state",
ExportJobState.EXPORTING);
+ overLimitJobMgr.unprotectAddJob(overLimitRunningJob);
+ overLimitJobMgr.removeOldExportJobs();
+
+ assertSameLabelIsStillUsed(overLimitJobMgr);
+ } finally {
+ Config.max_export_history_job_num = originalMaxHistoryJobNum;
+ }
+ }
+
+ @Test
+ public void testReplayUpdateForMissingExportJob() {
+ ExportMgr isolatedExportMgr = new ExportMgr();
+ ExportJob missingJob = makeExportJob(5001, "missing-job");
+
+ for (ExportJobState state : ExportJobState.values()) {
+ Assertions.assertDoesNotThrow(() ->
isolatedExportMgr.replayUpdateJobState(
+ new ExportJobStateTransfer(missingJob, state)));
+ }
+ }
+
+ private void assertSameLabelIsStillUsed(ExportMgr isolatedExportMgr)
throws Exception {
+ try {
+ isolatedExportMgr.addExportJobAndRegisterTask(makeExportJob(5001,
"same-label"));
+ Assertions.fail("Expected the running job to keep the label
reserved");
+ } catch (LabelAlreadyUsedException expected) {
+ // Expected: the running job still owns the label.
+ }
+ }
+
private ExportJob makeExportJob(long id, String label) {
ExportJob job1 = new ExportJob(id);
Deencapsulation.setField(job1, "label", label);
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]