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]
