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]

Reply via email to