goutamadwant opened a new issue, #12010: URL: https://github.com/apache/seatunnel/issues/12010
### Search before asking - [x] I searched existing issues and pull requests and found no report covering this submission-acknowledgement case. Related work such as #7693, #9532, and #11653 covers pending-job scheduling and failover, but not an interrupted queue insertion being acknowledged as a successful submission. ### What happened On the current `dev` branch, `CoordinatorService.submitJob` initializes the `JobMaster`, creates a `PendingJobInfo`, calls `pendingJobQueue.put(...)`, and then completes `jobSubmitFuture` successfully. `PeekBlockingQueue.put` catches `InterruptedException`, logs it, and returns `void`. The interrupt is not restored and the failure is not propagated to `CoordinatorService`. As a result, the submission path continues as if insertion succeeded: 1. `jobSubmitFuture.complete(null)` reports success to the caller. 2. The physical plan is updated to `PENDING`. 3. The coordinator logs that the job entered the pending queue. 4. The pending queue does not contain the job. This can occur during a coordinator shutdown or master transition. `clearCoordinatorService()` calls `executorService.shutdownNow()`, which can interrupt an in-flight submission worker immediately before the queue insertion. Although the backing `LinkedBlockingQueue` is unbounded, its `put` operation acquires an internal lock interruptibly and can therefore throw before inserting. Distributed job metadata is written before the insertion, so a later active master may restore the job. This report does not assume permanent job loss. The confirmed problem is that the original submission can be reported as successful even though the active coordinator did not accept the job into its pending queue, making execution dependent on later recovery. Relevant code: - [`PeekBlockingQueue.put`](https://github.com/apache/seatunnel/blob/b03cf1f1697ed84d0f57f416673cec40a075ef5a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/utils/PeekBlockingQueue.java#L59-L70) - [`CoordinatorService.submitJob`](https://github.com/apache/seatunnel/blob/b03cf1f1697ed84d0f57f416673cec40a075ef5a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/CoordinatorService.java#L1441-L1463) - [`CoordinatorService.clearCoordinatorService`](https://github.com/apache/seatunnel/blob/b03cf1f1697ed84d0f57f416673cec40a075ef5a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/CoordinatorService.java#L1280-L1330) ### Expected behavior and proposed contract A successful submission response should mean that the job was accepted into the pending scheduling path, or into another explicitly guaranteed recovery path. Before preparing a fix, please confirm which contract is intended: 1. **Fail interrupted or stale submission:** if interruption occurs before the job reaches an accepted scheduling or recovery boundary, complete `jobSubmitFuture` exceptionally with a retriable failure and preserve the thread interrupt status. Cleanup must interrupt the initialized `JobMaster` and remove partial state only when it is still owned by this submission, without deleting state already claimed or restored by a new active master. 2. **Make acceptance non-interruptible and lifecycle-safe:** because the queue is unbounded, use a non-interruptible insertion operation, but coordinate insertion with the coordinator activation/shutdown epoch so no job can be enqueued after `clearCoordinatorService()` has cleared the local queue. Acknowledge success only after the queue and job-id index are updated consistently, or after another explicitly documented durable-recovery boundary has been reached. `restoreJobFromMasterActiveSwitch()` also uses `pendingJobQueue.put(...)`. If the queue contract changes to propagate interruption, the restore path needs an explicit interruption and epoch policy as well. ### Reproduction The behavior is deterministic when the submission task is paused immediately before the real queue insertion and `clearCoordinatorService()` interrupts the coordinator executor: The injected queue deliberately widens the narrow production race window immediately before `PeekBlockingQueue.put` acquires its outer lock. It still uses the real coordinator executor and the production `clearCoordinatorService()` shutdown path. The interrupt comes from `executorService.shutdownNow()`, not from the test thread. ```java private static class InterruptOnShutdownPendingJobQueue extends PeekBlockingQueue<PendingJobInfo> { private final CountDownLatch putStarted = new CountDownLatch(1); private final CountDownLatch neverReleased = new CountDownLatch(1); private InterruptOnShutdownPendingJobQueue() { super(PendingJobInfo::getJobId); } @Override public void put(PendingJobInfo element) { putStarted.countDown(); try { neverReleased.await(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } super.put(element); } } InterruptOnShutdownPendingJobQueue queue = new InterruptOnShutdownPendingJobQueue(); ReflectionUtils.setField(coordinatorService, "pendingJobQueue", queue); PassiveCompletableFuture<Void> submitFuture = coordinatorService.submitJob(jobId, data, false); assertTrue(queue.putStarted.await(20, TimeUnit.SECONDS)); coordinatorService.clearCoordinatorService(); assertDoesNotThrow(submitFuture::join); assertFalse(queue.contains(jobId)); assertEquals(0, queue.size()); assertTrue(runningJobInfoIMap.containsKey(jobId)); assertEquals(JobStatus.PENDING, runningJobStateIMap.get(jobId)); ``` The focused `CoordinatorServiceTest` passed while asserting this inconsistent result on both JDK 8 and JDK 11: ```shell ./mvnw -pl seatunnel-engine/seatunnel-engine-server \ -Dskip.spotless=true -Dcheckstyle.skip=true -Dlicense.skip=true \ -Dtest=CoordinatorServiceTest#reproduceInterruptedPendingJobInsertionReportsSuccess test ``` ### Observed log ```log ERROR PeekBlockingQueue - Put element into queue failed. java.lang.InterruptedException INFO PhysicalPlan - Job interrupted_pending_job_insertion (...) turned from state CREATED to PENDING. INFO CoordinatorService - The submit job enter the pending queue ... ``` At that point, the submission future has completed normally while the pending queue has no entry for the job. The distributed `runningJobInfoIMap` entry still exists and `runningJobStateIMap` reports `PENDING`, which is why later coordinator recovery may be possible even though the original acceptance response was false. ### SeaTunnel Version Current `dev` at `b03cf1f1697ed84d0f57f416673cec40a075ef5a`. ### SeaTunnel Config ```conf The issue is not configuration-specific. The reproduction uses the existing seatunnel-engine-server test resource: batch_fake_to_console.conf. ``` ### Running Command ```shell ./mvnw -pl seatunnel-engine/seatunnel-engine-server \ -Dskip.spotless=true -Dcheckstyle.skip=true -Dlicense.skip=true \ -Dtest=CoordinatorServiceTest#reproduceInterruptedPendingJobInsertionReportsSuccess test ``` ### Error Exception ```log java.lang.InterruptedException at java.util.concurrent.locks.AbstractQueuedSynchronizer.acquireInterruptibly(...) at java.util.concurrent.LinkedBlockingQueue.put(...) at org.apache.seatunnel.engine.server.utils.PeekBlockingQueue.put(PeekBlockingQueue.java:62) at org.apache.seatunnel.engine.server.CoordinatorService.lambda$submitJob$11(CoordinatorService.java:1450) ``` ### Zeta or Flink or Spark Version Zeta engine from the current `dev` branch. ### Java or Scala Version Reproduced independently with JDK 8 and JDK 11. ### Are you willing to submit a PR? Yes, after maintainers confirm the intended interruption contract. ### Code of Conduct - [x] I agree to follow the project's Code of Conduct. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
