This is an automated email from the ASF dual-hosted git repository.
zhangshenghang pushed a commit to branch zsh-sync-st3517-pending-scheduler
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to
refs/heads/zsh-sync-st3517-pending-scheduler by this push:
new 1610a7cac1 [Fix][Zeta] Serialize master activation so pending jobs are
not dropped
1610a7cac1 is described below
commit 1610a7cac10050b2fcfc041b69fc9c4bf25072f4
Author: Shenghang <[email protected]>
AuthorDate: Thu Aug 6 23:05:46 2026 +0800
[Fix][Zeta] Serialize master activation so pending jobs are not dropped
CoordinatorServiceTest#testFailoverStopsOldPendingQueueAndNewCoordinatorCanSchedule
failed on CI with a 30s ConditionTimeoutException (CountDownLatch expected
0 but
was 1): the freshly activated coordinator never ran the job that was already
sitting in its pending queue.
Root cause: checkNewActiveMaster() was not mutually exclusive. Two threads
could
both observe isActive == false and both execute the activation block, so
pendingJobScheduleEpoch was bumped more than once with no step-down in
between.
The scheduler thread started by the first activation then found its own
epoch
stale once preApplyResources() returned and took the stale branch, which
interrupts the JobMaster and removes the PendingJobInfo it had already
reserved
from pendingJobQueue. Because no clearCoordinatorService() had run,
restoreAllRunningJobFromMasterNodeSwitch never rebuilt that entry, so the
job was
silently lost and run() was never invoked. In the failing test the two
callers are
the constructor-scheduled masterActiveListener tick and the reflective
checkNewActiveMaster() call, which the test cannot fully order.
Make checkNewActiveMaster() synchronized, on the same monitor already used
by
clearCoordinatorService(), so that "check ownership then flip isActive" is
atomic
and the schedule epoch can only change on a real activation or a real
step-down.
Add a regression test that drives four concurrent checkNewActiveMaster()
calls and
asserts the epoch is bumped exactly once and the queued job still runs
exactly
once. Without this fix it reports epoch 4 instead of 1. Also add the
comment on
isPendingJobSchedulerCurrent requested in review.
---
.../engine/server/CoordinatorService.java | 20 +++++-
.../engine/server/CoordinatorServiceTest.java | 75 ++++++++++++++++++++++
2 files changed, 94 insertions(+), 1 deletion(-)
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 95b3b0dad3..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
@@ -334,6 +334,12 @@ 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;
}
@@ -1219,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(
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 3438806054..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);