DanielLeens opened a new pull request, #12316:
URL: https://github.com/apache/seatunnel/pull/12316
## Summary
Checkpoint and savepoint barriers can be starved by a busy source for an
unbounded number of poll cycles, because the reader thread and the barrier
injector share one unfair intrinsic monitor and the reader re-enters it
immediately after every poll. This PR adds an explicit reader-to-injector
handoff in `SourceFlowLifeCycle` so a barrier is injected within at most one
`pollNext` call, for every source connector, without touching the public
`Collector#getCheckpointLock()` contract.
Observed on `dev`'s own `engine-v2-it` runs as
`BackpressureSlowSinkIT#testCheckpointsKeepCompletingUnderSustainedBackpressure`
failing at both `:184` (first checkpoint never completes inside 2 min; the
coordinator then hits the 100s `checkpoint.timeout`) and `:247` (`only observed
0`/`2` new checkpoints in the 90s window). Both JDK 8 and JDK 11 legs are
affected (dev runs `34801745423` and `34821874806`).
## Root cause (verified against current `dev` source)
Execution chain for one barrier under backpressure:
1. `CheckpointBarrierTriggerOperation#runInternal` hands
`SourceSeaTunnelTask#triggerBarrier` to the task group's async executor, which
calls `SourceFlowLifeCycle#triggerBarrier` and blocks on `synchronized
(collector.getCheckpointLock())`.
2. The reader thread is inside `SourceReader#pollNext`, which by contract
holds that same lock for the whole call. `FakeSourceReader#pollNext` emits up
to `MAX_ROWS_PER_POLL = 4096` rows per call, and each
`SeaTunnelSourceCollector#collect` -> `sendRecordToNext` blocks on the bounded
intermediate queue (`ArrayBlockingQueue`, capacity 2048) while the sink drains
at ~500 rows/s. One poll therefore holds the lock for ~8 s.
3. `SourceFlowLifeCycle#collect` releases the lock when `pollNext` returns,
executes `Thread.sleep(0L)`, and re-enters `pollNext`. HotSpot monitors do not
hand off to the parked waiter on exit; the exiting thread that immediately
re-acquires wins nearly every time. `Thread.sleep(0L)` only yields the CPU and
does not change that race.
4. The barrier is injected only when the scheduler happens to favour the
injector at a release point, i.e. after a random integer number of ~8 s polls.
The dev log for job `1150683685079482369` shows exactly that shape: checkpoint
durations of 20 s, 34 s, 8.5 s and 36 s against a no-starvation baseline of ~6
s (2048-row intermediate queue plus the 1024-row `MultiTableSinkWriter` queue
at 500 rows/s).
#12313 makes the E2E fixture avoid the race by sizing FakeSource splits so
every split completes in one poll and its explicit `Thread.sleep(1000L)` runs
outside the lock; that PR itself notes the engine-side unfair handoff as
pre-existing and out of its scope. This PR is that engine-side fix. Any
production source whose `pollNext` is long relative to the checkpoint interval
(large batches, slow downstream, JDBC/file readers with big fetch sizes) is
exposed to the same starvation; #11489 already worked around it inside
`FakeSourceReader` by batching, which is why the lock hold is ~8 s rather than
the whole split.
## Fix
- New package-private `SourceCheckpointLockHandoff`
(`seatunnel-engine-server`, `task.flow`): an injector announces itself before
contending for the lock and withdraws in `finally`; the reader, at the point
where it holds no lock, waits (1 ms slices, accounted as source idle time when
observability is enabled) until no injector is announced before starting the
next poll.
- `SourceFlowLifeCycle#triggerBarrier` wraps its `synchronized` block with
`injectorArriving()` / `injectorFinished()`.
- `SourceFlowLifeCycle#collect` replaces the ineffective `Thread.sleep(0L)`
after a non-empty poll with `checkpointLockHandoff.awaitInjectors()`.
- No change to `Collector`, `SourceReader`, any connector, checkpoint
serialization, restore, or barrier ordering. The barrier still travels through
the same lock and the same `sendRecordToNext` call, after `addState` and `ack`,
exactly as before.
### Semantics and safety
- Old behaviour: barrier injection latency is unbounded and depends on
scheduler luck. New behaviour: bounded by the poll in flight. Data order, ack
order and state snapshot content are unchanged.
- Deadlock-free by construction: the reader waits only while holding no
lock; an injector never waits for the reader; an injector blocked on the full
intermediate queue while forwarding the barrier still makes progress because
the sink keeps draining. The `finally` guarantees a failed injection (for
example `snapshotState` throwing) cannot leave the reader parked.
- Hot-path cost when no barrier is pending: one volatile read per non-empty
poll (`awaitInjectors` returns 0 immediately).
- Empty polls are unaffected (they already sleep 100 ms outside the lock);
`prepareClose` and schema-change phases are unaffected (they run after the
handoff point and do not hold the lock).
- Task cancellation still interrupts the reader thread; the wait propagates
`InterruptedException` like the previous sleep did.
- The periodic `FlushSignal` from `onTimerTick` also takes the checkpoint
lock; it is intentionally left out of the handoff (it is not a barrier and is
not latency-critical), keeping the change to the barrier path only.
## Tests
- New `SourceCheckpointLockHandoffTest` (pure unit test, no cluster):
- no injector pending: `awaitInjectors` returns immediately (no added
latency on the hot path);
- pending state tracks multiple injectors independently (checkpoint and
savepoint triggered back to back);
- a reader re-acquiring the lock in a tight loop yields to an injector
within at most the poll in flight (measured in poll cycles, not wall-clock);
- a parked reader resumes when the injector finishes, including after a
failed injection;
- a parked reader honours interrupt (task cancellation).
- Existing E2E coverage that exercises the barrier-vs-busy-source path:
`BackpressureSlowSinkIT`, `SavepointBusySourceBarrierIT`, plus the whole
`engine-v2-it` suite.
- Local: `./mvnw spotless:apply -pl seatunnel-engine/seatunnel-engine-server
-nsu -Dmaven.gitcommitid.skip=true` only. Compilation, unit tests and E2E are
verified by this PR's GitHub CI, per this initiative's no-local-build policy.
## Files
-
`seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/flow/SourceFlowLifeCycle.java`
(handoff wiring, comments updated)
-
`seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/flow/SourceCheckpointLockHandoff.java`
(new)
-
`seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/task/flow/SourceCheckpointLockHandoffTest.java`
(new)
🤖 Generated with [Claude Code](https://claude.com/claude-code)
--
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]