Vivek1106-04 commented on PR #12511: URL: https://github.com/apache/seatunnel/pull/12511#issuecomment-5982000557
@DanielLeens thanks for the thorough review. All eight points are addressed in `9143d479c` and `a691fb195`. The branch is also rebased onto `dev` `af80a6704`, so it is 0 behind and picks up #12444 for the PayPal flake. **Issue 1 (gate has no test that fails without it).** New `testTriggerAfterInFlightCheckpointEndsButBeforeSavepointIsCreatedRearms` uses your lock-holding idea, with a real drain thread. The test thread takes the coordinator `lock` and drops `pendingCounter` to 0. It then waits until the savepoint caller has woken and is `BLOCKED` on the lock, and fires a trigger of each type; the lock is reentrant, so those calls run. Each trigger must create nothing and re-arm. After the lock is released, exactly one savepoint is created. With `|| savepointDraining` removed it fails: `a CHECKPOINT_TYPE trigger created a checkpoint ahead of the savepoint ==> expected: <0> but was: <1>`. @Rangsh ported the matching #12493 test, thanks. I went with this version instead because it covers the same window without setting the gate by reflection, so no gate is left without a request behind it. **Issue 2 (re-interrupt on a pool worker).** The flag is now re-set only when the caller is not a `ForkJoinWorkerThread`. `testInterruptedSavepointDrainFailsTheRequestAndResumesTriggering` asserts that a caller on its own thread keeps the flag. New `testInterruptedSavepointDrainOnPoolWorkerClearsTheInterruptFlag` runs the drain on a `ForkJoinPool` and asserts that the flag is clear when `startSavepoint()` returns. It fails with the unconditional interrupt (`expected: <false> but was: <true>`). **Issue 3 (cancel during the restore wait ends FAILED).** Your inference was right, and I confirmed it by running it. New `CheckpointErrorRestoreEndTest#testCancelDuringPipelineRestoreWaitEndsTheJobCanceled` fails the first checkpoint, waits until the pipeline is reset for restore (30 s wait in its conf), and cancels: - previous head `91001c42c`: `expected: <CANCELED> but was: <FAILED>` - `dev`: passes. The log shows `PhysicalPlan` going `RUNNING -> CANCELING`, then `cancelPipeline` waiting on the SubPlan monitor until the wait ends. The pipeline restarts (`CREATED -> SCHEDULED -> DEPLOYING -> RUNNING`), is cancelled straight away, and the job ends `CANCELED`. I kept `dev`'s outcome. If the job is `CANCELING` when the restore is abandoned, the pipeline ends `CANCELED`; otherwise it ends in the state it failed with and keeps its error. The savepoint case is unchanged, and the cancelled job no longer restarts a pipeline only to cancel it. **Issue 4 (root cause and point-in-time check).** You were right that `forceStop` is not a no-op. I took a log trace of the restore-wait `SavePointTest` on `dev`'s `SubPlan`. During the wait, `forceStop` moves the enumerator and committer `CREATED -> CANCELED`, and both log `state process is not start`. Their vertex state process was stopped, so `taskFuture` is never completed. Only the source task's future completes. The pipeline therefore never counts all tasks ended. 2 s later the restart goes `SCHEDULED -> DEPLOYING -> RUNNING` and deploys nothing, because `DEPLOYING` starts only `CREATED` tasks. The pipeline then sits in `RUNNING`. So the lost stop comes from `PhysicalVertex.forceStop` on a task whose state process is not running, not from callbacks queued behind the monitor. The monitor does explain why `cancelPipeline` waits for the restart in the Issue 3 scenario. The Javadoc and the PR description now describe this mechanism. On the residual window: a stop decided after the re-check but before deploy hits the same `forceStop` gap. A second `isNeedRestore()` check in `SCHEDULED` would narrow the window, but it is still check-then-act against a stop path that takes no SubPlan lock, so it cannot close it. The real fix is for `forceStop` to complete the future of a task whose state process is not running. That changes a stop path shared by every job, so I would rather do it as a separate issue and PR than grow this one. Happy to open it if you agree. **Issue 5 (wall-clock SavePointTest).** The conf now pins `job.retry.interval.seconds = 3`. The fixed `Thread.sleep(2000L)` is replaced by an await on the job being `RUNNING` with every pipeline's coordinator reporting all tasks ready, so the savepoint cannot hit `TASK_NOT_ALL_READY_WHEN_SAVEPOINT`. **Issue 6 (futures completed under the lock, no log).** `failDrainingSavepoint` clears the gate under the lock and completes the request after leaving it. `cleanPendingCheckpoint` does the same: it captures the draining request under the lock and fails it after the synchronized block. That is safe for the reset window you pointed out, because `restoreCoordinator()` clears `shutdown` only after `cleanPendingCheckpoint` returns, and by then the request is done. Interrupted and cut-short drains now `LOG.warn` with job id, pipeline id and close reason. **Issue 7.** `drainAndTriggerSavepoint` has a Javadoc stating the loop invariant. `savepointPendingCheckpoint` is now `volatile`, and its comment says that only tests read it. **Issue 8.** #12493 is closed in favour of this PR. I would keep the `SubPlan` change here: the two commits stay separate, and without it any correct drain fix turns `SavePointTest` red. **Coverage gap from 2.2.** The reset test is now parameterized over `CHECKPOINT_COORDINATOR_RESET`, `PIPELINE_END`, `CHECKPOINT_COORDINATOR_COMPLETED` and `CHECKPOINT_INSIDE_ERROR`, and each must fail the request with `CHECKPOINT_COORDINATOR_SHUTDOWN`. Local results on JDK 17 at `a691fb195`: `CheckpointCoordinatorTest` 25/25, `SavePointTest` 6 run + 1 existing `@Disabled`, `CheckpointErrorRestoreEndTest` 2/2. `spotless:apply` makes no changes and `-DskipTests verify` passes. The red evidence for each new test is in the updated PR description. CI is running on the new head. -- 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]
