agarwalrahul2702 opened a new pull request, #11489:
URL: https://github.com/apache/seatunnel/pull/11489

   # [Fix][Zeta] Fix stop-with-savepoint hang in DOING_SAVEPOINT
   
   Fixes #11473
   
   ### Root cause
   
   Two defects combined to produce the reported hang:
   
   1. **Savepoint barriers could never be injected while the source was 
emitting a large split.**
      `SourceFlowLifeCycle.triggerBarrier` must acquire the checkpoint lock to 
inject a
      checkpoint/savepoint barrier into the source. `FakeSourceReader.pollNext` 
emitted the
      **entire split** in a single call while holding that same lock. With the 
issue's config
      (`row.num = 100000000`, `split.num = 1`) the lock is held for the whole 
100M-row split —
      effectively forever. A thread dump taken during the hang shows the 
barrier thread
      `BLOCKED` on the checkpoint lock whose owner is the source task thread 
busy inside
      `FakeDataGenerator` (the sink meanwhile idle-polling an empty queue — 
this is not
      backpressure). Consequently every checkpoint and savepoint expires at
      `checkpoint.timeout`, which also explains `SinkCommittedCount = 0` in the 
issue report.
   
   2. **A failed savepoint left the job wedged in `DOING_SAVEPOINT` forever.**
      `JobMaster.savePoint()` sets `DOING_SAVEPOINT`, triggers savepoint 
barriers, and waits.
      When the savepoint checkpoint expires, the checkpoint error handling 
restores the
      pipeline and the job keeps running — but nothing ever reverted the job 
status. The job
      stayed in `DOING_SAVEPOINT` indefinitely with `errorMsg = null` (exactly 
the reported
      symptom), while `/finished-jobs/SAVEPOINT_DONE` never listed it and slots 
stayed
      occupied until a manual normal stop.
   
   ### Solution
   
   - `FakeSourceReader` now emits at most 4096 rows per `pollNext` call and 
requeues the
     remainder of the split at the head of the split deque, releasing the 
checkpoint lock
     between batches so barriers can be injected. The requeued split carries 
the remaining
     row count, so a snapshot taken between batches captures the 
not-yet-emitted rows
     (state semantics preserved for restore). Split-read-interval pacing and 
the idle sleep
     are preserved at split granularity; user-configured custom rows (`rows` 
option) are
     emitted in one batch as before.
   - `JobMaster.savePoint()` now reverts the job state from `DOING_SAVEPOINT` 
back to
     `RUNNING` (new `PhysicalPlan#savepointFailed`) when the savepoint does not 
complete, so
     the failure is surfaced to the caller instead of wedging the job. No 
timeout or forced
     cancellation is introduced; the savepoint completion path itself is 
unchanged.
   
   ### Reproduction (before the fix, on current dev)
   
   In-process Zeta cluster (`AbstractSeaTunnelServerTest`) running the issue's 
exact
   FakeSource → Console streaming config:
   
   - Job reaches `RUNNING`, savepoint requested → status turns 
`DOING_SAVEPOINT`, the
     savepoint pending checkpoint is never created, the in-flight checkpoint 
expires after
     `checkpoint.timeout`, the pipeline restores, and the job stays 
`DOING_SAVEPOINT`
     indefinitely; the stop request eventually fails with
     `CheckpointException: CheckpointCoordinator shutdown` while the status 
remains wedged.
   - A baseline run of the same config **without** any savepoint shows every 
periodic
     checkpoint expiring too (barrier starvation is independent of the 
savepoint request).
   - A normal stop (`isStopWithSavePoint=false`) cleans the job up, as reported.
   
   ### Regression test
   
   `SavePointBusySourceTest` (engine-server) proves, against a real in-process 
Zeta cluster:
   
   1. stop-with-savepoint **completes** under a busy source emitting a huge 
split;
   2. the job passes `DOING_SAVEPOINT` → `SAVEPOINT_DONE`;
   3. the stop request (the same `CoordinatorService#savePoint` call the REST
      `/stop-job` handler uses) returns instead of hanging;
   4. all slots are released (`SlotService` worker profile shows 0 assigned 
slots);
   5. a normal stop while the source is busy still works (`CANCELED`, slots 
released);
   6. a savepoint that cannot complete (sink delays `prepareCommit` past
      `checkpoint.timeout`) reports an `ExecutionException` to the caller and 
the job
      reverts to `RUNNING` instead of sticking in `DOING_SAVEPOINT`, then 
reaches a
      terminal state on stop with slots released.
   
   ### Verification (local, JDK 11 / Temurin 11.0.31, macOS arm64)
   
   | Command | Result |
   |---|---|
   | `./mvnw -pl seatunnel-engine/seatunnel-engine-server 
-Dtest=SavePointBusySourceTest test` | Tests run: 3, Failures: 0, Errors: 0 |
   | `./mvnw -pl seatunnel-engine/seatunnel-engine-server -Dtest=SavePointTest 
test` | Tests run: 6, Failures: 0, Errors: 0, Skipped: 1 (pre-existing 
`@Disabled`) |
   | `./mvnw -pl seatunnel-connectors-v2/connector-fake test` | Tests run: 12, 
Failures: 0, Errors: 0 |
   | `./mvnw -pl seatunnel-engine/seatunnel-engine-server 
-Dtest=CheckpointTimeOutTest,CheckpointBarrierTriggerErrorTest,CheckpointPlanTest,CheckpointErrorRestoreEndTest,CheckpointManagerTest,CheckpointCoordinatorTest,CheckpointSerializeTest,JobMasterTest
 test` | Tests run: 24, Failures: 0, Errors: 0, Skipped: 2 (pre-existing) |
   | `./mvnw spotless:apply` (changed modules) | BUILD SUCCESS (formatting 
applied) |
   | `./mvnw -DskipTests verify` (full repo, 284 modules) | BUILD SUCCESS 
(55:46 min) |
   
   ### Checklist
   
   - [x] Code changed follows the project code style (spotless applied)
   - [x] New files carry the ASF license header
   - [x] Regression test added
   - [x] No user-facing config options changed; no docs impact
   - [x] No new dependencies
   


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