This is an automated email from the ASF dual-hosted git repository.
davidzollo pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new f13b610eff [Fix][Zeta] Avoid duplicate pending job scheduling after
failover (#11653)
f13b610eff is described below
commit f13b610effa5769c79e1903969ee4869525fb9a9
Author: Jast <[email protected]>
AuthorDate: Sun Aug 23 16:16:50 2026 +0800
[Fix][Zeta] Avoid duplicate pending job scheduling after failover (#11653)
Co-authored-by: David Zollo <[email protected]>
---
.../engine/server/CoordinatorService.java | 278 +++++++++++++++------
.../engine/server/utils/PeekBlockingQueue.java | 63 ++++-
.../engine/server/CoordinatorServiceTest.java | 237 ++++++++++++++++++
3 files changed, 498 insertions(+), 80 deletions(-)
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/CoordinatorService.java
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/CoordinatorService.java
index 15c76c8835..328d245095 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/CoordinatorService.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/CoordinatorService.java
@@ -122,6 +122,7 @@ import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLong;
+import java.util.concurrent.locks.ReentrantLock;
import java.util.stream.Collectors;
import static
org.apache.seatunnel.api.options.EnvCommonOptions.CHECKPOINT_INTERVAL;
@@ -185,6 +186,25 @@ public class CoordinatorService {
private final PeekBlockingQueue<PendingJobInfo> pendingJobQueue =
new PeekBlockingQueue<>(PendingJobInfo::getJobId);
+ // Fair lock that serializes "peek the head job and reserve it for
evaluation" across all
+ // scheduler generations. Holding this lock around peek+reserve guarantees
that at most one
+ // thread holds a given job id in schedulingPendingJobIds at any time, so
two overlapping
+ // scheduler threads (e.g. a stale one still inside preApplyResources and
a fresh one started
+ // after a master flap) cannot independently reserve the same job.
+ private final ReentrantLock pendingJobScheduleLock = new
ReentrantLock(true);
+ // Monotonically increasing generation counter for scheduler runs. Each
scheduler thread
+ // snapshots this value at submission time and exits when it is bumped.
Bumped inside
+ // checkNewActiveMaster and clearCoordinatorService to invalidate any
in-flight scheduler.
+ private final AtomicLong pendingJobScheduleEpoch = new AtomicLong();
+ // JobMasters currently inside jobMaster.preApplyResources(). Interrupted
en masse by
+ // clearCoordinatorService so they release any pending RPC and stop
touching shared state.
+ private final Set<JobMaster> schedulingJobMasters =
ConcurrentHashMap.newKeySet();
+ // Job ids currently reserved by a scheduler thread for evaluation under
+ // pendingJobScheduleLock. reservePendingJobInfo adds;
releasePendingJobInfo removes. Cleared
+ // by clearCoordinatorService so a freshly-restored scheduler can
re-reserve jobs after a
+ // master flap.
+ private final Set<Long> schedulingPendingJobIds =
ConcurrentHashMap.newKeySet();
+
/**
* IMap key is {@link PipelineLocation}
*
@@ -286,21 +306,27 @@ public class CoordinatorService {
*/
private void startPendingJobScheduleThread() {
logger.info("Start pending job schedule thread");
+ long scheduleEpoch = pendingJobScheduleEpoch.get();
Runnable pendingJobScheduleTask =
() -> {
Thread.currentThread().setName("pending-job-schedule-runner");
- while (isActive) {
+ while (isPendingJobSchedulerCurrent(scheduleEpoch)) {
try {
- pendingJobSchedule();
- } catch (InterruptedException interrupted) {
- throw new RuntimeException(interrupted);
+ pendingJobSchedule(scheduleEpoch);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ return;
} catch (Throwable e) {
+ if (!isPendingJobSchedulerCurrent(scheduleEpoch)) {
+ return;
+ }
logger.severe("Error in pending job schedule
thread", e);
try {
Thread.sleep(3000L);
} catch (InterruptedException ex) {
logger.severe("Pending job schedule thread
interrupted", ex);
Thread.currentThread().interrupt();
+ return;
}
}
}
@@ -308,6 +334,50 @@ public class CoordinatorService {
executorService.submit(pendingJobScheduleTask);
}
+ /**
+ * A scheduler generation is "current" only when this node is still the
active master AND no
+ * newer generation has been started since this thread was submitted.
Every scheduler-loop
+ * iteration and every stale-branch decision hinges on this predicate:
once it turns false the
+ * thread must stop touching shared scheduling state, because a newer
generation now owns it.
+ */
+ private boolean isPendingJobSchedulerCurrent(long scheduleEpoch) {
+ return isActive && pendingJobScheduleEpoch.get() == scheduleEpoch;
+ }
+
+ /**
+ * Reserves the next pending job for evaluation under
pendingJobScheduleLock. The returned
+ * PendingJobInfo has its job id added to schedulingPendingJobIds,
blocking any other scheduler
+ * thread from reserving the same job until {@link #releasePendingJobInfo}
is called. Returns
+ * null only if the queue is empty.
+ */
+ private PendingJobInfo reservePendingJobInfo() throws InterruptedException
{
+ pendingJobScheduleLock.lockInterruptibly();
+ try {
+ PendingJobInfo pendingJobInfo =
+ pendingJobQueue.peekBlocking(
+ jobInfo ->
!schedulingPendingJobIds.contains(jobInfo.getJobId()));
+ if (Objects.isNull(pendingJobInfo)) {
+ return null;
+ }
+ schedulingPendingJobIds.add(pendingJobInfo.getJobId());
+ return pendingJobInfo;
+ } finally {
+ pendingJobScheduleLock.unlock();
+ }
+ }
+
+ /**
+ * Releases a previously reserved pending job and wakes any other
scheduler thread that may be
+ * blocked inside reservePendingJobInfo waiting for a non-empty queue.
Idempotent: safe to call
+ * when the job was never reserved.
+ */
+ private void releasePendingJobInfo(PendingJobInfo pendingJobInfo) {
+ if (Objects.nonNull(pendingJobInfo)
+ && schedulingPendingJobIds.remove(pendingJobInfo.getJobId())) {
+ pendingJobQueue.release();
+ }
+ }
+
/**
* Attempts to advance the head job in the pending queue into execution.
*
@@ -317,10 +387,19 @@ public class CoordinatorService {
* {@link JobMaster#run()}. Restored jobs rehydrate pipeline execution
state before the new
* master takes over execution.
*
+ * <p>If {@code scheduleEpoch} no longer matches {@link
#pendingJobScheduleEpoch} when this
+ * method finishes evaluating resources (because the node stepped down or
master ownership
+ * transferred), the job's {@link JobMaster} is interrupted, its entry is
removed from the
+ * pending queue, and the method returns without invoking {@code run()}. A
new scheduler thread
+ * started for the current epoch will pick up (or rebuild) the job from
{@code
+ * runningJobInfoIMap} instead of re-running the poisoned instance.
+ *
* @throws InterruptedException if queue waiting or retry sleep is
interrupted
+ * @throws Throwable any unchecked failure from {@code preApplyResources};
the reservation is
+ * released before propagation so a future scheduler iteration can
retry.
*/
- private void pendingJobSchedule() throws InterruptedException {
- PendingJobInfo pendingJobInfo = pendingJobQueue.peekBlocking();
+ private void pendingJobSchedule(long scheduleEpoch) throws Throwable {
+ PendingJobInfo pendingJobInfo = reservePendingJobInfo();
if (Objects.isNull(pendingJobInfo)) {
// This situation almost never happens because pendingJobSchedule
is single-threaded
logger.warning("The peek job info is null");
@@ -337,71 +416,88 @@ public class CoordinatorService {
String.format(
"Start calculating whether pending task resources are
enough: %s", jobId));
- boolean preApplyResources = jobMaster.preApplyResources();
- if (!preApplyResources) {
- try {
- PendingJobDiagnostic diagnostic =
- PendingDiagnosticsCollector.collectJobDiagnostic(
- pendingJobInfo, Collections.emptyMap(),
getResourceManager());
- pendingJobInfo.recordSnapshot(diagnostic);
- } catch (Exception e) {
- logger.warning(
- String.format(
- "Collect pending diagnostic for job %s failed:
%s",
- jobId, ExceptionUtils.getMessage(e)));
+ boolean preApplyResources;
+ schedulingJobMasters.add(jobMaster);
+ try {
+ preApplyResources = jobMaster.preApplyResources();
+ } catch (Throwable e) {
+ releasePendingJobInfo(pendingJobInfo);
+ throw e;
+ } finally {
+ schedulingJobMasters.remove(jobMaster);
+ }
+ try {
+ if (!isPendingJobSchedulerCurrent(scheduleEpoch)) {
+ jobMaster.interrupt();
+ // Drop the poisoned PendingJobInfo from the queue: its
JobMaster has been
+ // interrupted (which permanently completes
jobMasterCompleteFuture exceptionally),
+ // so the next-epoch scheduler must not re-dispatch it.
Restore on the next
+ // activation will rebuild a fresh PendingJobInfo from
runningJobInfoIMap.
+ pendingJobQueue.remove(pendingJobInfo);
+ return;
}
- logger.info(
- String.format(
- "Current strategy is %s, and resources is not
enough, skipping this schedule, JobID: %s",
- scheduleStrategy, jobId));
- if (isWaitStrategy) {
+ if (!preApplyResources) {
try {
- Thread.sleep(3000);
- } catch (InterruptedException e) {
- logger.severe(ExceptionUtils.getMessage(e));
+ PendingJobDiagnostic diagnostic =
+ PendingDiagnosticsCollector.collectJobDiagnostic(
+ pendingJobInfo, Collections.emptyMap(),
getResourceManager());
+ pendingJobInfo.recordSnapshot(diagnostic);
+ } catch (Exception e) {
+ logger.warning(
+ String.format(
+ "Collect pending diagnostic for job %s
failed: %s",
+ jobId, ExceptionUtils.getMessage(e)));
}
- return;
- } else {
- completeFailJob(jobMaster);
- queueRemove(jobMaster);
- return;
- }
- }
- logger.info(String.format("Resources enough, start running: %s",
jobId));
- // When deleting jobmaster from pendingJobQueue, make sure that there
is a corresponding
- // jobMaster in the runningJobMasterMap
- runningJobMasterMap.put(jobId, jobMaster);
- final PendingJobInfo finalPendingJobInfo = pendingJobQueue.take();
- final JobMaster finalJobMaster = finalPendingJobInfo.getJobMaster();
- PendingSourceState pendingSourceState =
finalPendingJobInfo.getPendingSourceState();
- MDCExecutorService mdcExecutorService = MDCTracer.tracing(jobId,
executorService);
- mdcExecutorService.submit(
- () -> {
+ logger.info(
+ String.format(
+ "Current strategy is %s, and resources is not
enough, skipping this schedule, JobID: %s",
+ scheduleStrategy, jobId));
+ if (isWaitStrategy) {
try {
- String jobFullName =
finalJobMaster.getPhysicalPlan().getJobFullName();
- JobStatus jobStatus = (JobStatus)
runningJobStateIMap.get(jobId);
- if (pendingSourceState == PendingSourceState.RESTORE
- && !jobStatus.isEndState()) {
- finalJobMaster
- .getPhysicalPlan()
- .getPipelineList()
- .forEach(SubPlan::restorePipelineState);
- }
- logger.info(
- String.format(
- "The %s %s is in %s state, restore
pipeline and take over this job running",
- pendingSourceState, jobFullName,
jobStatus));
- finalJobMaster.run();
- } finally {
- if (jobMasterCompletedSuccessfully(finalJobMaster,
pendingSourceState)) {
- runningJobMasterMap.remove(jobId);
- }
+ Thread.sleep(3000);
+ } catch (InterruptedException e) {
+ logger.severe(ExceptionUtils.getMessage(e));
}
- });
- }
-
- private void queueRemove(JobMaster jobMaster) {
- pendingJobQueue.removeById(jobMaster.getJobId());
+ return;
+ } else {
+ completeFailJob(jobMaster);
+ pendingJobQueue.remove(pendingJobInfo);
+ return;
+ }
+ }
+ logger.info(String.format("Resources enough, start running: %s",
jobId));
+ // When deleting jobmaster from pendingJobQueue, make sure that
there is a corresponding
+ // jobMaster in the runningJobMasterMap
+ runningJobMasterMap.put(jobId, jobMaster);
+ pendingJobQueue.remove(pendingJobInfo);
+ PendingSourceState pendingSourceState =
pendingJobInfo.getPendingSourceState();
+ MDCExecutorService mdcExecutorService = MDCTracer.tracing(jobId,
executorService);
+ mdcExecutorService.submit(
+ () -> {
+ try {
+ String jobFullName =
jobMaster.getPhysicalPlan().getJobFullName();
+ JobStatus jobStatus = (JobStatus)
runningJobStateIMap.get(jobId);
+ if (pendingSourceState ==
PendingSourceState.RESTORE
+ && !jobStatus.isEndState()) {
+ jobMaster
+ .getPhysicalPlan()
+ .getPipelineList()
+
.forEach(SubPlan::restorePipelineState);
+ }
+ logger.info(
+ String.format(
+ "The %s %s is in %s state, restore
pipeline and take over this job running",
+ pendingSourceState, jobFullName,
jobStatus));
+ jobMaster.run();
+ } finally {
+ if (jobMasterCompletedSuccessfully(jobMaster,
pendingSourceState)) {
+ runningJobMasterMap.remove(jobId);
+ }
+ }
+ });
+ } finally {
+ releasePendingJobInfo(pendingJobInfo);
+ }
}
/**
@@ -1129,8 +1225,20 @@ public class CoordinatorService {
* <p>When this node becomes the active master, the coordinator
initializes distributed services
* and triggers job restore. When it loses master ownership, local
coordinator state is torn
* down. Initialization failures are cleaned up locally and retried by
later polling cycles.
+ *
+ * <p>Synchronized on the same monitor as {@link
#clearCoordinatorService()} so that the "check
+ * ownership then flip {@code isActive}" sequence is atomic. Two threads
running this method
+ * concurrently could otherwise both observe {@code isActive == false} and
both perform the
+ * activation block, bumping {@link #pendingJobScheduleEpoch} twice
without any step-down in
+ * between. The scheduler thread started by the first activation would
then see its own epoch as
+ * stale once {@code preApplyResources} returns and discard the pending
job it had already
+ * reserved (interrupt its JobMaster and drop it from {@code
pendingJobQueue}) even though this
+ * node never actually lost master ownership. Because no {@code
clearCoordinatorService()} ran,
+ * {@code restoreAllRunningJobFromMasterNodeSwitch} is never triggered to
rebuild that entry, so
+ * the job would be silently lost. Epoch changes must therefore only ever
come from a real
+ * activation or a real step-down.
*/
- private void checkNewActiveMaster() {
+ private synchronized void checkNewActiveMaster() {
try {
if (!isActive && this.seaTunnelServer.isMasterNode()) {
logger.info(
@@ -1139,6 +1247,7 @@ public class CoordinatorService {
this.executorService = createCoordinatorExecutor();
}
initCoordinatorService();
+ pendingJobScheduleEpoch.incrementAndGet();
isActive = true;
startPendingJobScheduleThread();
seaTunnelServer.startRealtimeMetricsService(this);
@@ -1174,19 +1283,32 @@ public class CoordinatorService {
if (!coordinatorServiceCleared.compareAndSet(false, true)) {
return;
}
+ pendingJobScheduleEpoch.incrementAndGet();
+ schedulingJobMasters.forEach(JobMaster::interrupt);
+ schedulingJobMasters.clear();
+ schedulingPendingJobIds.clear();
+ pendingJobQueue.release();
// interrupt all JobMaster
runningJobMasterMap.values().forEach(JobMaster::interrupt);
- if (isWaitStrategy) {
- pendingJobQueue
- .getJobIdMap()
- .values()
- .forEach(
- pendingJobInfo -> {
- JobMaster jobMaster =
pendingJobInfo.getJobMaster();
- jobMaster.interrupt();
- });
- pendingJobQueue.clear();
- }
+ // Interrupt and discard every JobMaster currently sitting in
pendingJobQueue. This is
+ // intentionally unconditional (not gated on isWaitStrategy): every
entry in
+ // pendingJobQueue corresponds to a JobInfo already persisted in
runningJobInfoIMap, so
+ // clearCoordinatorService can safely drop the local queue copy.
Restore on the next
+ // activation will rebuild PendingJobInfo entries from
runningJobInfoIMap via
+ // restoreAllRunningJobFromMasterNodeSwitch, producing a single fresh
JobMaster per job
+ // id. Leaving the stale PendingJobInfo in the queue (the pre-fix
behavior for
+ // REJECT strategy) caused an interrupted JobMaster's
permanently-poisoned
+ // jobMasterCompleteFuture to be picked by the next-epoch scheduler
and led to duplicate
+ // physicalPlan.startJob() invocations for the same job id.
+ pendingJobQueue
+ .getJobIdMap()
+ .values()
+ .forEach(
+ pendingJobInfo -> {
+ JobMaster jobMaster =
pendingJobInfo.getJobMaster();
+ jobMaster.interrupt();
+ });
+ pendingJobQueue.clear();
executorService.shutdownNow();
runningJobMasterMap.clear();
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/utils/PeekBlockingQueue.java
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/utils/PeekBlockingQueue.java
index 699f5bca08..b866d6fbdb 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/utils/PeekBlockingQueue.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/utils/PeekBlockingQueue.java
@@ -29,6 +29,7 @@ import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
import java.util.function.Function;
+import java.util.function.Predicate;
/**
* PeekBlockingQueue implements blocking when peeking. Queues like
BlockingQueue only support
@@ -76,18 +77,57 @@ public class PeekBlockingQueue<E> {
return element;
}
+ /**
+ * Wakes up any thread currently blocked inside {@link
#peekBlocking(Predicate)} waiting for the
+ * queue to become non-empty. No-op when the queue is empty; a thread
waiting on an
+ * actually-empty queue still relies on the caller's {@code
shutdownNow}/{@code interrupt} to
+ * unblock, so do not remove the {@code shutdownNow} call in {@code
+ * CoordinatorService.clearCoordinatorService} assuming this method is
sufficient on its own.
+ */
+ public void release() {
+ lock.lock();
+ try {
+ if (!queue.isEmpty()) {
+ notEmpty.signalAll();
+ }
+ } finally {
+ lock.unlock();
+ }
+ }
+
public E peekBlocking() throws InterruptedException {
+ return peekBlocking(element -> true);
+ }
+
+ /**
+ * Blocks until the queue contains an element satisfying {@code predicate}
(in FIFO order) and
+ * returns it without removing. The caller is responsible for removing the
element via {@link
+ * #remove(Object)} (or {@link #take()}) once the element is no longer
needed so the next waiter
+ * can make progress on a different element.
+ */
+ public E peekBlocking(Predicate<E> predicate) throws InterruptedException {
lock.lock();
try {
- while (queue.peek() == null) {
+ E element = findFirst(predicate);
+ while (element == null) {
notEmpty.await();
+ element = findFirst(predicate);
}
- return queue.peek();
+ return element;
} finally {
lock.unlock();
}
}
+ private E findFirst(Predicate<E> predicate) {
+ for (E element : queue) {
+ if (predicate.test(element)) {
+ return element;
+ }
+ }
+ return null;
+ }
+
public Integer size() {
lock.lock();
try {
@@ -124,6 +164,25 @@ public class PeekBlockingQueue<E> {
}
}
+ /**
+ * Removes a specific element by identity from both the underlying queue
and the id-to-element
+ * map. Uses {@link Map#remove(Object, Object)} on the id map so a newer
entry that happens to
+ * share the same id (e.g. a freshly restored {@code PendingJobInfo} that
has overwritten the id
+ * pointer) is not accidentally evicted alongside the intended removal.
+ */
+ public boolean remove(E element) {
+ lock.lock();
+ try {
+ boolean removed = queue.remove(element);
+ if (removed) {
+ jobIdMap.remove(idExtractor.apply(element), element);
+ }
+ return removed;
+ } finally {
+ lock.unlock();
+ }
+ }
+
public boolean contains(Long jobId) {
return jobIdMap.containsKey(jobId);
}
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java
index 286ba67fa6..7e385e38f7 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java
@@ -87,6 +87,7 @@ import java.util.Map;
import java.util.Set;
import java.util.concurrent.CompletionException;
import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.CyclicBarrier;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
@@ -95,6 +96,7 @@ import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicLong;
import static
org.apache.seatunnel.engine.core.classloader.DefaultClassLoaderService.SKIP_CHECK_JAR;
import static org.awaitility.Awaitility.await;
@@ -418,6 +420,79 @@ public class CoordinatorServiceTest {
}
}
+ /**
+ * Concurrent {@code checkNewActiveMaster()} calls must activate the
coordinator exactly once.
+ *
+ * <p>Regression test for a race in which two threads both observed {@code
isActive == false}
+ * and both ran the activation block, bumping the pending-job schedule
epoch twice without any
+ * step-down in between. The scheduler thread started by the first
activation then saw its own
+ * epoch as stale after {@code preApplyResources()} returned and discarded
the job it had
+ * already reserved — interrupting its JobMaster and dropping it from
{@code pendingJobQueue} —
+ * even though the node never lost master ownership. Because no {@code
+ * clearCoordinatorService()} ran, nothing rebuilt that entry, so the job
was silently lost and
+ * never ran. This surfaced as a 30s {@code ConditionTimeoutException} in
{@link
+ * #testFailoverStopsOldPendingQueueAndNewCoordinatorCanSchedule()}, whose
reflective {@code
+ * checkNewActiveMaster()} call could interleave with the
constructor-scheduled
+ * masterActiveListener tick.
+ */
+ @Test
+ void testConcurrentMasterActivationActivatesOnceAndKeepsPendingJob()
throws Exception {
+ // Start as non-master so the constructor-scheduled
masterActiveListener cannot activate
+ // before we stop it; that keeps this test's two threads the only
activation drivers.
+ AtomicBoolean masterFlag = new AtomicBoolean(false);
+ SeaTunnelServer server = Mockito.mock(SeaTunnelServer.class);
+ Mockito.when(server.isMasterNode()).thenAnswer(invocation ->
masterFlag.get());
+ CoordinatorService coordinatorService =
newMockCoordinatorService(server);
+ int activationThreads = 4;
+ ExecutorService activationPool =
Executors.newFixedThreadPool(activationThreads);
+ try {
+ CountDownLatch runLatch = new CountDownLatch(1);
+ JobMaster pendingJob = enqueueMockPendingJob(coordinatorService,
80001L, runLatch);
+ long epochBefore =
getPendingJobScheduleEpoch(coordinatorService).get();
+
+ masterFlag.set(true);
+ CyclicBarrier startTogether = new CyclicBarrier(activationThreads);
+ List<Future<?>> activations = new ArrayList<>();
+ for (int i = 0; i < activationThreads; i++) {
+ activations.add(
+ activationPool.submit(
+ () -> {
+ startTogether.await(30, TimeUnit.SECONDS);
+
invokeCheckNewActiveMaster(coordinatorService);
+ return null;
+ }));
+ }
+ for (Future<?> activation : activations) {
+ activation.get(60, TimeUnit.SECONDS);
+ }
+
+ await().atMost(30, TimeUnit.SECONDS)
+ .untilAsserted(
+ () ->
Assertions.assertTrue(coordinatorService.isCoordinatorActive()));
+ Assertions.assertEquals(
+ epochBefore + 1,
+ getPendingJobScheduleEpoch(coordinatorService).get(),
+ "concurrent checkNewActiveMaster calls must bump the
schedule epoch only once");
+ await().atMost(30, TimeUnit.SECONDS)
+ .untilAsserted(() -> Assertions.assertEquals(0L,
runLatch.getCount()));
+ Mockito.verify(pendingJob, Mockito.times(1)).run();
+ Mockito.verify(pendingJob, Mockito.never()).interrupt();
+
Assertions.assertFalse(coordinatorService.getPendingJobQueue().contains(80001L));
+ } finally {
+ activationPool.shutdownNow();
+ shutdownCoordinatorIfRunning(coordinatorService);
+ }
+ }
+
+ private AtomicLong getPendingJobScheduleEpoch(CoordinatorService
coordinatorService) {
+ return ReflectionUtils.getField(coordinatorService,
"pendingJobScheduleEpoch")
+ .map(AtomicLong.class::cast)
+ .orElseThrow(
+ () ->
+ new AssertionError(
+ "Failed to get pendingJobScheduleEpoch
by reflection"));
+ }
+
@Test
void testCheckNewActiveMasterIsIdempotentWhenAlreadyActive() throws
Exception {
AtomicBoolean masterFlag = new AtomicBoolean(true);
@@ -444,6 +519,149 @@ public class CoordinatorServiceTest {
}
}
+ @Test
+ void testPendingJobSchedulerIgnoresJobReservedByPreviousScheduler() throws
Exception {
+ SeaTunnelServer server = Mockito.mock(SeaTunnelServer.class);
+ CoordinatorService coordinatorService =
newMockCoordinatorService(server);
+ CountDownLatch firstScheduleStarted = new CountDownLatch(1);
+ CountDownLatch allowFirstScheduleToFinish = new CountDownLatch(1);
+ try {
+ CountDownLatch runLatch = new CountDownLatch(1);
+ JobMaster jobMaster = enqueueMockPendingJob(coordinatorService,
70001L, runLatch);
+ AtomicBoolean firstSchedule = new AtomicBoolean(true);
+ Mockito.when(jobMaster.preApplyResources())
+ .thenAnswer(
+ invocation -> {
+ if (firstSchedule.compareAndSet(true, false)) {
+ firstScheduleStarted.countDown();
+ allowFirstScheduleToFinish.await();
+ }
+ return true;
+ });
+
+ ReflectionUtils.setField(coordinatorService, "isActive", true);
+ invokePendingJobScheduler(coordinatorService);
+ Assertions.assertTrue(firstScheduleStarted.await(5,
TimeUnit.SECONDS));
+
+ ReflectionUtils.setField(coordinatorService, "isActive", false);
+ ReflectionUtils.setField(coordinatorService, "isActive", true);
+ invokePendingJobScheduler(coordinatorService);
+
+ await().during(500, TimeUnit.MILLISECONDS)
+ .atMost(1, TimeUnit.SECONDS)
+ .untilAsserted(
+ () -> Mockito.verify(jobMaster,
Mockito.times(1)).preApplyResources());
+
+ allowFirstScheduleToFinish.countDown();
+ await().atMost(5, TimeUnit.SECONDS)
+ .untilAsserted(
+ () ->
+ Assertions.assertFalse(
+ coordinatorService
+ .getPendingJobQueue()
+ .contains(70001L)));
+ Mockito.verify(jobMaster, Mockito.atMostOnce()).run();
+ } finally {
+ allowFirstScheduleToFinish.countDown();
+ shutdownCoordinatorIfRunning(coordinatorService);
+ }
+ }
+
+ @Test
+ void
testPendingJobSchedulerCanAdvanceNextJobWhenPreviousResourceCheckBlocks()
+ throws Exception {
+ SeaTunnelServer server = Mockito.mock(SeaTunnelServer.class);
+ CoordinatorService coordinatorService =
newMockCoordinatorService(server);
+ CountDownLatch firstScheduleStarted = new CountDownLatch(1);
+ CountDownLatch allowFirstScheduleToFinish = new CountDownLatch(1);
+ try {
+ JobMaster blockedJobMaster =
+ enqueueMockPendingJob(coordinatorService, 80001L, new
CountDownLatch(1));
+ Mockito.when(blockedJobMaster.preApplyResources())
+ .thenAnswer(
+ invocation -> {
+ firstScheduleStarted.countDown();
+ while (!allowFirstScheduleToFinish.await(
+ 100, TimeUnit.MILLISECONDS)) {
+ // Simulate a resource request that does
not react immediately
+ // to master demotion or scheduler
replacement.
+ }
+ return true;
+ });
+
+ ReflectionUtils.setField(coordinatorService, "isActive", true);
+ invokePendingJobScheduler(coordinatorService);
+ Assertions.assertTrue(firstScheduleStarted.await(5,
TimeUnit.SECONDS));
+
+ JobMaster secondJobMaster =
+ enqueueMockPendingJob(coordinatorService, 80002L, new
CountDownLatch(1));
+ invokePendingJobScheduler(coordinatorService);
+
+ await().atMost(5, TimeUnit.SECONDS)
+ .untilAsserted(
+ () -> {
+ Mockito.verify(secondJobMaster,
Mockito.times(1))
+ .preApplyResources();
+ Assertions.assertFalse(
+
coordinatorService.getPendingJobQueue().contains(80002L));
+ });
+ Mockito.verify(blockedJobMaster,
Mockito.times(1)).preApplyResources();
+ } finally {
+ allowFirstScheduleToFinish.countDown();
+ shutdownCoordinatorIfRunning(coordinatorService);
+ }
+ }
+
+ @Test
+ void testClearCoordinatorServiceDropsPendingJobsUnderRejectStrategy()
throws Exception {
+ // Regression test for the duplicate-dispatch gap on the default
REJECT schedule strategy:
+ // when the coordinator is cleared while a scheduler is blocked inside
preApplyResources,
+ // the interrupted PendingJobInfo must NOT survive in pendingJobQueue,
otherwise a
+ // restored scheduler (post same-node master flap-back) could
re-dispatch the poisoned
+ // JobMaster after a fresh PendingJobInfo under the same job id has
been enqueued.
+ SeaTunnelServer server = Mockito.mock(SeaTunnelServer.class);
+ EngineConfig engineConfig = new EngineConfig();
+ engineConfig.setScheduleStrategy(ScheduleStrategy.REJECT);
+ CoordinatorService coordinatorService =
newMockCoordinatorService(server, engineConfig);
+ CountDownLatch allowFirstScheduleToFinish = new CountDownLatch(1);
+ try {
+ JobMaster blockedJobMaster =
+ enqueueMockPendingJob(coordinatorService, 90001L, new
CountDownLatch(1));
+ Mockito.when(blockedJobMaster.preApplyResources())
+ .thenAnswer(
+ invocation -> {
+ allowFirstScheduleToFinish.await();
+ return true;
+ });
+
+ ReflectionUtils.setField(coordinatorService, "isActive", true);
+ invokePendingJobScheduler(coordinatorService);
+
+ // Wait until the scheduler thread is parked inside
preApplyResources().
+ await().atMost(5, TimeUnit.SECONDS)
+ .untilAsserted(
+ () ->
+ Mockito.verify(blockedJobMaster,
Mockito.atLeastOnce())
+ .preApplyResources());
+
+ // Simulate a master step-down. The blocked JobMaster is
interrupted; the
+ // PendingJobInfo must be dropped from the queue so a later
restore cannot
+ // re-dispatch the poisoned instance.
+ invokeClearCoordinatorService(coordinatorService);
+
+ await().atMost(5, TimeUnit.SECONDS)
+ .untilAsserted(
+ () -> {
+ Assertions.assertFalse(
+
coordinatorService.getPendingJobQueue().contains(90001L));
+ Mockito.verify(blockedJobMaster,
Mockito.atLeastOnce()).interrupt();
+ });
+ } finally {
+ allowFirstScheduleToFinish.countDown();
+ shutdownCoordinatorIfRunning(coordinatorService);
+ }
+ }
+
@Test
void testPendingJobWithInsufficientResourceRespectsWaitStrategy() throws
Exception {
AtomicBoolean masterFlag = new AtomicBoolean(true);
@@ -548,6 +766,12 @@ public class CoordinatorServiceTest {
.ifPresent(ScheduledExecutorService::shutdownNow);
}
+ private void invokePendingJobScheduler(CoordinatorService
coordinatorService) throws Exception {
+ Method method =
CoordinatorService.class.getDeclaredMethod("startPendingJobScheduleThread");
+ method.setAccessible(true);
+ method.invoke(coordinatorService);
+ }
+
private JobMaster enqueueMockPendingJob(
CoordinatorService coordinatorService, long jobId, CountDownLatch
runLatch) {
return enqueueMockPendingJob(coordinatorService, jobId, runLatch,
true);
@@ -594,6 +818,12 @@ public class CoordinatorServiceTest {
"Failed to get coordinator
executorService by reflection"));
}
+ private void shutdownCoordinatorIfRunning(CoordinatorService
coordinatorService) {
+ if (!getCoordinatorExecutor(coordinatorService).isShutdown()) {
+ coordinatorService.shutdown();
+ }
+ }
+
@Test
public void
testSeaTunnelEngineRetryableExceptionOperationCanBeRetryByHazelcast() {
@@ -1227,6 +1457,13 @@ public class CoordinatorServiceTest {
method.invoke(coordinatorService);
}
+ private void invokeClearCoordinatorService(CoordinatorService
coordinatorService)
+ throws Exception {
+ Method method =
CoordinatorService.class.getDeclaredMethod("clearCoordinatorService");
+ method.setAccessible(true);
+ method.invoke(coordinatorService);
+ }
+
@Test
public void testClearCoordinatorService() {
JobInformation jobInformation =