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);

Reply via email to