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]

Reply via email to