Skip to content

Commit 36c2966

Browse files
committed
[Fix][Zeta] Ignore stale cleanup fence for restored job generation
1 parent 7278866 commit 36c2966

2 files changed

Lines changed: 53 additions & 1 deletion

File tree

  • seatunnel-engine/seatunnel-engine-server/src

seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/master/JobMaster.java

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -989,7 +989,10 @@ public boolean isStatePersistenceAllowed() {
989989
}
990990
IMap<Long, JobCleanupRecord> pendingJobCleanupIMap =
991991
nodeEngine.getHazelcastInstance().getMap(Constant.IMAP_PENDING_JOB_CLEANUP);
992-
return !pendingJobCleanupIMap.containsKey(jobId);
992+
JobCleanupRecord cleanupRecord = pendingJobCleanupIMap.get(jobId);
993+
return cleanupRecord == null
994+
|| !Objects.equals(
995+
cleanupRecord.getOwnerInitializationTimestamp(), initializationTimestamp);
993996
}
994997

995998
/**

seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/master/JobMasterTest.java

Lines changed: 49 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -390,6 +390,55 @@ void testCleanupFenceBlocksMissingStateWriterBetweenSnapshotAndRegistration() th
390390
}
391391
}
392392

393+
@Test
394+
void testOldGenerationCleanupDoesNotFenceCurrentGeneration() throws Exception {
395+
long jobId = instance.getFlakeIdGenerator(Constant.SEATUNNEL_ID_GENERATOR_NAME).newId();
396+
long oldInitializationTimestamp = 1L;
397+
long currentInitializationTimestamp = 2L;
398+
JobMaster jobMaster =
399+
newJobMaster(
400+
jobId,
401+
"batch_fake_to_console.conf",
402+
"test_restore_cleanup_generation_fence",
403+
false);
404+
IMap<Long, JobCleanupRecord> pendingJobCleanupIMap =
405+
nodeEngine.getHazelcastInstance().getMap(Constant.IMAP_PENDING_JOB_CLEANUP);
406+
407+
try {
408+
runningJobInfoIMap.put(jobId, new JobInfo(currentInitializationTimestamp, null));
409+
jobMaster.init(currentInitializationTimestamp, false);
410+
411+
pendingJobCleanupIMap.put(
412+
jobId,
413+
new JobCleanupRecord(
414+
oldInitializationTimestamp,
415+
JobStatus.FINISHED,
416+
Collections.singleton(jobId),
417+
Collections.singleton(jobId),
418+
System.currentTimeMillis()));
419+
Assertions.assertTrue(
420+
jobMaster.isStatePersistenceAllowed(),
421+
"A cleanup record from an older generation must not fence a restored job with the same job id");
422+
423+
pendingJobCleanupIMap.put(
424+
jobId,
425+
new JobCleanupRecord(
426+
currentInitializationTimestamp,
427+
JobStatus.FINISHED,
428+
Collections.singleton(jobId),
429+
Collections.singleton(jobId),
430+
System.currentTimeMillis()));
431+
Assertions.assertFalse(
432+
jobMaster.isStatePersistenceAllowed(),
433+
"The cleanup record for the current generation must still fence late writers");
434+
} finally {
435+
pendingJobCleanupIMap.remove(jobId);
436+
runningJobInfoIMap.remove(jobId);
437+
runningJobStateIMap.remove(jobId);
438+
runningJobStateTimestampsIMap.remove(jobId);
439+
}
440+
}
441+
393442
private void assertCloseIdleTask(JobMaster jobMaster) {
394443
SlotService slotService = server.getSlotService();
395444
long jobId = jobMaster.getJobId();

0 commit comments

Comments
 (0)