Vivek1106-04 opened a new pull request, #12511: URL: https://github.com/apache/seatunnel/pull/12511
### Purpose of this pull request Part of #12441. It fixes the same-coordinator part of the issue, and two bugs that any correct fix for it exposes. #12493 does not cover those two bugs (see https://github.com/apache/seatunnel/pull/12493#issuecomment-5862502658). **1. Savepoint drain holds the coordinator lock (the issue itself)** `CheckpointCoordinator.startSavepoint()` sleep-polled `pendingCounter` inside `synchronized (lock)`. Every `tryTriggerPendingCheckpoint()` on that coordinator therefore blocked for the whole in-flight checkpoint. `LastCheckpointNotifyOperation` can also reach that path on a Hazelcast operation thread. - The drain now runs outside the lock. - A `savepointDraining` gate and one shared `savepointRequest` future are installed together under the lock. Concurrent `startSavepoint()` calls get that same future for the whole drain. - While the gate is set, every non-savepoint trigger (general, completed-point, schema-change) re-arms, as it already does behind a pending checkpoint. The gate check sits next to `pendingCounter > 0` in `tryTriggerPendingCheckpoint()`. It is not in `createPendingCheckpoint()`, because that method cannot re-arm. - Interrupt, shutdown, completion, reset and create failure all clear the gate under the lock and fail the shared future, so ordinary triggering is never left suppressed. **2. Failure reason of a drain cut short** `JobMaster.isSavepointStartPreconditionException` treats a savepoint as never started, and retryable, only for `CHECKPOINT_COORDINATOR_SHUTDOWN` and `TASK_NOT_ALL_READY_WHEN_SAVEPOINT`. When `cleanPendingCheckpoint` ends a drain, whatever the close reason (`CHECKPOINT_COORDINATOR_RESET`, `PIPELINE_END`, `CHECKPOINT_COORDINATOR_COMPLETED`, or an error reason), it now always reports `CHECKPOINT_COORDINATOR_SHUTDOWN`. Passing the close reason through makes the job master treat a pipeline ending during the drain as a real savepoint failure, even though no savepoint checkpoint was ever created, and breaks `SavePointTest`. **3. SubPlan restore race (job stuck in `DOING_SAVEPOINT`)** This bug already exists on `dev`, but the lock-held drain hides it because it delays `cleanPendingCheckpoint`. Once the drain is off the lock, `SavePointTest#testSavePointButJobGoingToFail` fails consistently in the full test run. A pipeline that fails while its job is in `DOING_SAVEPOINT` is reset and waits `pipelineRestoreIntervalSeconds` before restarting. `JobMaster.savepointFailed()` calls `forceStopPipeline()` during that wait, but finds only reset tasks and stops nothing, so the pipeline restarts and the job never leaves `DOING_SAVEPOINT`. `SubPlan` now re-checks `jobMaster.isNeedRestore()` after the wait. If a stop was decided meanwhile, the pipeline goes back to the terminal state it failed with and keeps its error. **Not in this PR:** the cross-pipeline regression (part 2 of the issue) needs the shared dispatch pool from #12165. It is written and red/green on top of #12165, and will follow once #12165 is merged, so #12441 stays open until then. The #12442 trigger-failure policy is untouched. ### Does this PR introduce _any_ user-facing change? No. Savepoint and checkpoint behavior is unchanged except that it no longer blocks or hangs in the cases above. No config options or defaults change. ### How was this patch tested? New tests in `CheckpointCoordinatorTest`. The drain is test-owned: the test holds `pendingCounter > 0` until its assertions finish, with no wall-clock sleeps. - `testTriggerDuringSavepointDrainReturnsAndRearmsWithoutCreatingACheckpoint`: a trigger of each type (`CHECKPOINT_TYPE`, `COMPLETED_POINT_TYPE`, `SCHEMA_CHANGE_BEFORE_POINT_TYPE`, `SCHEMA_CHANGE_AFTER_POINT_TYPE`) during a drain returns within 10 s, creates no checkpoint, and re-arms. Once the drain ends, exactly one savepoint is created. - `testConcurrentSavepointsDuringDrainShareOneSavepoint`: a second request returns at once with the same outcome, and only one checkpoint id is allocated. - `testInterruptedSavepointDrainFailsTheRequestAndResumesTriggering` - `testCoordinatorResetDuringSavepointDrainFailsTheRequestAndResumesTriggering`: the reset does not wait for the drain, and the request fails with `CHECKPOINT_COORDINATOR_SHUTDOWN`. New test in `SavePointTest`: - `testSavePointFailureDuringPipelineRestoreWaitEndsTheJob` with `stream_two_pipelines_savepoint_fails_during_restore.conf`. Sink A fails the savepoint at 4 s and sink B finishes its part at 5 s, so the stop lands inside the 3 s restore wait. Red on `dev` (`146a1b5c5`), with the production changes reverted and the tests kept: | Test | Result on `dev` | |---|---| | `testTriggerDuringSavepointDrainReturnsAndRearmsWithoutCreatingACheckpoint` | `a CHECKPOINT_TYPE trigger blocked behind the savepoint drain ==> Unexpected exception thrown: java.util.concurrent.TimeoutException` (10 s bound) | | `testConcurrentSavepointsDuringDrainShareOneSavepoint` | `a second savepoint request blocked behind the first one's drain ==> Unexpected exception thrown: java.util.concurrent.TimeoutException` | | `testInterruptedSavepointDrainFailsTheRequestAndResumesTriggering` | `ExecutionException: java.lang.InterruptedException: sleep interrupted` (`@SneakyThrows` rethrows instead of failing the request) | | `testCoordinatorResetDuringSavepointDrainFailsTheRequestAndResumesTriggering` | passes on `dev`; kept to pin down the `CHECKPOINT_COORDINATOR_SHUTDOWN` reason for the new drain | | `SavePointTest#testSavePointFailureDuringPipelineRestoreWaitEndsTheJob` | `ConditionTimeoutException: ... expected: <FAILED> but was: <DOING_SAVEPOINT> within 1 minutes.` | Green on this branch: ``` Tests run: 20, Failures: 0, Errors: 0, Skipped: 0 - in ...checkpoint.CheckpointCoordinatorTest Tests run: 7, Failures: 0, Errors: 0, Skipped: 1 - in ...checkpoint.SavePointTest (the skip is an existing @Disabled test) ``` Local runs on JDK 17. `./mvnw spotless:apply` shows no changes. ### Check list * [ ] If any new Jar binary package adding in your PR, please add License Notice according [New License Guide](https://github.com/apache/seatunnel/blob/dev/docs/en/developer/new-license.md) * [ ] If necessary, please update the documentation to describe the new feature. https://github.com/apache/seatunnel/tree/dev/docs * [ ] If necessary, please update `incompatible-changes.md` to describe the incompatibility caused by this PR. -- 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]
