SEPURI-SAI-KRISHNA commented on PR #12290: URL: https://github.com/apache/seatunnel/pull/12290#issuecomment-5770968922
Thanks, and to answer directly: **#12316** is the one to review for the backpressure failure. #12313 is worth reading first as context, but it is not the fix. They are not two competing implementations, they are the same defect at two levels, both by @DanielLeens. **#12316** is the engine-side root cause. `FakeSourceReader#pollNext` emits rows while holding the checkpoint lock that `SourceFlowLifeCycle#triggerBarrier` needs. When a split exceeds `MAX_ROWS_PER_POLL` the reader sets `splitInProgress` and skips its own `Thread.sleep(1000L)` at `FakeSourceReader.java:173`, so it re-acquires the lock back to back with no deterministic release point and barrier injection starves. #12316 adds `SourceCheckpointLockHandoff` and replaces the `Thread.sleep(0L)` between polls with `awaitInjectors()`, so the reader yields when an injector is pending. That is in `SourceFlowLifeCycle`, so it applies to every source, not just `FakeSource`. **#12313** changes only the test: the config on dev is `row.num = 10000000` with `split.num = 1`, one split far above the 4096 per-poll cap, and #12313 resizes it to `row.num = 2000000` / `split.num = 500` so every split is 4000 rows and the reader takes its real sleep between splits. Its Java diff is entirely javadoc, zero code lines. It makes this one test deterministic without touching the starvation itself, so if #12316 lands, #12313 becomes a test tweak rather than the fix. One correction to my own earlier framing while I am here. I described the `BackpressureSlowSinkIT` failure as hitting the sustained-window assertion. It does not. In the run I read it failed at the first-checkpoint gate, `completed >= 1` never becoming true inside 2 minutes while the job was healthy and RUNNING, which is the starvation above rather than a slow window. A small note if you review #12313 as well: the comment it adds says the reader's sleep gives barrier injection a contention-free window "once per second". With 4000 rows per split and the sink at `write_delay_ms = 2`, so 500 rows per second, a split takes 8 seconds, so the window is once per 8 seconds. The conclusion still holds against a 15 second `checkpoint.interval`, but the stated margin is 8x larger than the real one, and that margin is the whole justification for the chosen numbers. On this PR: agreed, leave it parked and rebase once the engine fixes land, then rerun `engine-v2-it`. Thanks for taking the cancellation path on #12311. -- 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]
