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

   ## Summary
   
   Follow-up to #12107. 
`BackpressureSlowSinkIT#testCheckpointsKeepCompletingUnderSustainedBackpressure`
 is failing on `dev` itself, making `dev`'s `engine-v2-it` red:
   
   - dev run `34821874806` (JDK 11): `expected at least 3 additional 
checkpoints ... only observed 0 (samples=[1,1,1,1,1,1,1,1,1])`
   - dev run `34801745423` (JDK 8): `only observed 2`
   
   No `src/main` commit between 09-04 and 09-14 touches the checkpoint lock, 
the intermediate queue, or the barrier path, so this is not a product 
regression — the test's own fixture makes barrier injection non-deterministic.
   
   ## Root cause (verified against current `dev` source)
   
   The previous fixture 
(`stream_fast_fakesource_to_slow_inmemory_backpressure.conf`) used `row.num = 
10000000`, `split.num = 1` — a single 10M-row split.
   
   - `FakeSourceReader#pollNext` emits at most `MAX_ROWS_PER_POLL = 4096` rows 
per call while holding the checkpoint lock 
(`seatunnel-connectors-v2/connector-fake/.../FakeSourceReader.java:47,100-149`) 
— the same lock `SourceFlowLifeCycle#triggerBarrier` acquires to inject a 
checkpoint/savepoint barrier 
(`seatunnel-engine/.../SourceFlowLifeCycle.java:481`).
   - Because the split had far more than 4096 rows remaining, every `pollNext` 
requeued the remainder and set `splitInProgress = true` 
(`FakeSourceReader.java:142-149`). That skips both the split-read-interval gate 
(`FakeSourceReader.java:96-98`) and the reader's own inter-poll 
`Thread.sleep(1000L)` (`FakeSourceReader.java:172-174`) — the only place that 
sleep runs fully outside the checkpoint lock.
   - The reader therefore re-entered the `synchronized` block back-to-back, 
poll after poll, with no deterministic release point for the barrier thread to 
win the intrinsic-lock race. Each poll holds the lock for ~8s at the sink's 
throttled ~500 rows/sec drain rate (`write_delay_ms = 2` in 
`InMemorySinkWriter#write`), so under unfair contention `triggerBarrier` can be 
starved for many consecutive polls — well past the 15s `checkpoint.interval` — 
exactly matching the CI logs.
   
   ## Fix (config + comments only — no assertion or bound touched)
   
   - `row.num = 2000000`, `split.num = 500`. 
`FakeSourceSplitEnumerator#discoverySplits` computes `splitRowNum = 
ceil(row.num / split.num)` (`FakeSourceSplitEnumerator.java:120`), so every 
split is exactly 4000 rows — at or under `MAX_ROWS_PER_POLL = 4096`.
   - Every split therefore completes in a single `pollNext` call 
(`remainingRowNum = 0`), so `splitInProgress` never gets set and the reader 
always takes its real `Thread.sleep(1000L)` fully outside the checkpoint lock 
between splits — giving `triggerBarrier` a guaranteed, contention-free 
one-second window after every split.
   - Backpressure stays genuinely sustained: the intermediate queue capacity is 
2048 (`TaskGroupWithIntermediateBlockingQueue.java:44`); at ~500 rows/sec, the 
1s gap between splits drains only ~500 rows (~25%) before the next split 
refills it, so the queue never empties and `emitBlockedNs` stays meaningful.
   - 2,000,000 rows at 500 rows/sec would take ~66 minutes to exhaust if fed 
continuously — far outlasting the test's 90s `BACKPRESSURE_WINDOW_MS`.
   - Updated the conf header/source comments and the `BackpressureSlowSinkIT` 
class Javadoc to document this split-sizing/lock-release coupling. No test 
assertion, Awaitility bound, `BACKPRESSURE_WINDOW_MS`, or 
`MIN_NEW_COMPLETED_CHECKPOINTS` changed — `emitBlockedNs > 0` still requires 
real, measured backpressure.
   
   ## Note for a separate discussion
   
   The underlying unfair intrinsic-lock handoff — a large-batch reader that 
keeps re-acquiring the checkpoint lock back-to-back can delay barrier injection 
without bound — is a pre-existing engine property. It is not addressed in this 
PR and is worth a separate discussion/issue.
   
   ## Test plan
   
   - Test-only change (config values + comments/Javadoc); no `src/main` or 
`pom.xml` touched.
   - `./mvnw spotless:apply -pl 
seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base -nsu 
-Dmaven.gitcommitid.skip=true` run locally, `BUILD SUCCESS`.
   - Verification is via this PR's GitHub CI (`engine-v2-it`), per this 
initiative's no-local-build policy for this task.
   
   🤖 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