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]