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]

Reply via email to