DanielLeens commented on PR #12316:
URL: https://github.com/apache/seatunnel/pull/12316#issuecomment-5678693760

   *Posting this as a plain issue comment rather than a `gh pr review` — GitHub 
does not let an author submit a formal review on their own PR. I reviewed this 
the same way I'd review anyone else's `seatunnel-engine` change: full diff, 
full call-chain trace, no benefit of the doubt for authorship.*
   
   # What Problem Does This PR Solve?
   
   - **User pain point**: Under sustained backpressure, a checkpoint or 
savepoint barrier injected via `SourceFlowLifeCycle#triggerBarrier` can be 
starved for a long, effectively unbounded number of reader poll cycles. The 
reader thread holds `collector.getCheckpointLock()` while emitting each record 
inside `SourceReader#pollNext`, releases it, and — because Java's intrinsic 
monitors are unfair — almost always wins the race to reacquire it for the next 
poll before a parked barrier-injector thread gets a turn. The previous 
mitigation, `Thread.sleep(0L)` after a non-empty poll, only yields the CPU; it 
does not change who wins the monitor race, so it did not bound the starvation. 
This was observed as `BackpressureSlowSinkIT` timing out on `dev` (per my prior 
root-cause note on this flake).
   - **Fix approach**: Add `SourceCheckpointLockHandoff`, a small 
cooperative-handoff primitive: a barrier injector calls `injectorArriving()` 
before contending for the lock and `injectorFinished()` in a `finally` block 
after releasing it; the reader thread, at the point between polls where it 
holds no lock, calls `awaitInjectors()` and spin-waits (1ms `Thread.sleep` 
slices) until no injector is registered before starting its next poll.
   - **One-sentence summary**: Turns the previously-unfair, 
potentially-unbounded race for the checkpoint lock into an explicit handoff 
that guarantees a barrier injector acquires the lock within one poll cycle of 
announcing itself, instead of possibly losing the race indefinitely.
   
   # 1. Code Change Review
   
   ## 1.1 Core Logic Analysis
   
   **Core changes**:
   - New file `SourceCheckpointLockHandoff.java` (102 lines): an `AtomicInteger 
pendingInjectors` counter with 
`injectorArriving()`/`injectorFinished()`/`hasPendingInjector()`/`awaitInjectors()`.
   - `SourceFlowLifeCycle.java`: `collect()` now calls 
`checkpointLockHandoff.awaitInjectors()` (instead of `Thread.sleep(0L)`) right 
after a non-empty poll, outside the checkpoint lock; `triggerBarrier(Barrier)` 
now wraps its existing `synchronized (collector.getCheckpointLock())` block in 
`injectorArriving()` / `try { ... } finally { injectorFinished(); }`.
   
   Before (`collect()`):
   ```java
   Thread.sleep(0L);
   ```
   After:
   ```java
   long handoffWaitNs = checkpointLockHandoff.awaitInjectors();
   if (metricsEnabled) {
       sourceIdleNs.inc(handoffWaitNs);
   }
   ```
   Before (`triggerBarrier`):
   ```java
   synchronized (collector.getCheckpointLock()) {
       ...
       collector.sendRecordToNext(new Record<>(barrier));
   }
   ```
   After:
   ```java
   checkpointLockHandoff.injectorArriving();
   try {
       synchronized (collector.getCheckpointLock()) {
           ...
           collector.sendRecordToNext(new Record<>(barrier));
       }
   } finally {
       checkpointLockHandoff.injectorFinished();
   }
   ```
   
   **Key findings**:
   - The normal path reaches this code on every source poll cycle for every 
running source subtask: `awaitInjectors()` is called unconditionally after 
every non-empty `pollNext()`, and `injectorArriving()`/`injectorFinished()` 
wrap every `triggerBarrier` call, i.e. every checkpoint and savepoint. This is 
core scheduling-path code, not a rare recovery branch.
   - `awaitInjectors()`'s own contract is respected at its call site: its 
Javadoc requires the caller to hold no lock, and 
`SourceFlowLifeCycle.collect()` calls it strictly after 
`reader.pollNext(collector)` has already returned (i.e. after any lock the 
reader took inside `pollNext` has already been released) — I traced this and it 
is correct, no deadlock is introduced.
   - `injectorFinished()` is unconditionally reached via `finally`, even when 
`barrier.prepareClose(...)`, `reader.snapshotState(...)`, 
`runningTask.addState(...)`, `runningTask.ack(...)`, or 
`collector.sendRecordToNext(...)` throws — so a failed injection cannot 
permanently park the reader. Verified there is exactly one 
`injectorArriving()`/`injectorFinished()` pair, with no early return between 
them.
   - `AtomicInteger` correctly supports the documented "more than one injector" 
case (a checkpoint and a savepoint barrier close together) — 
increments/decrements are safe under concurrent callers, and 
`hasPendingInjector()` correctly reflects "at least one still pending" until 
every announced injector has withdrawn.
   
   **In-depth correctness analysis — the residual bound this fix actually 
provides**: I want to be precise about what "bounded by a single poll" means, 
because I think the class Javadoc could be misread as a stronger guarantee than 
it actually is, and I'd rather flag my own doc wording than let it stand 
ambiguous.
   
   `SourceReader#pollNext` is documented as "generate the **next batch** of 
records" (not one record), and at least one connector — 
`JdbcSourceReader#pollNext` 
(`seatunnel-connectors-v2/connector-jdbc/.../JdbcSourceReader.java:67-95`) — 
wraps its **entire per-split read loop** (`while (!inputFormat.reachedEnd()) { 
output.collect(seaTunnelRow); }`) inside one `synchronized 
(output.getCheckpointLock())` block, i.e. it takes the lock once for the whole 
split rather than once per row. For that connector shape, `awaitInjectors()` 
cannot help *during* a single `pollNext()` call — the reader thread is inside 
`reader.pollNext(collector)`, not at the handoff call site, for the full 
duration of that split. What this fix actually bounds is the 
previously-unbounded compounding across **multiple consecutive poll cycles** 
(the reader winning the race poll after poll after poll, which is exactly what 
I traced as the mechanism behind the `BackpressureSlowSinkIT` flake — a test 
source that re
 turns from `pollNext()` frequently). For a connector whose single `pollNext()` 
call can itself run for minutes (e.g. a very large JDBC split under a slow 
downstream), checkpoint-barrier latency for that specific poll is not improved 
by this PR — it is unchanged from before, and was already a pre-existing 
property of holding the lock for the whole split, not something this PR 
regresses.
   - This means the fix is fully correct and effective for the flake it targets 
and for any connector that calls `collect()` once (or a few times) per 
`pollNext()` invocation, which covers most non-batch-oriented sources. It does 
not, and does not claim in the code, extend to eliminating lock hold time 
*within* a single large-batch poll.
   
   ## 1.2 Compatibility Impact
   
   **Fully compatible.** No config option, API, checkpoint/savepoint state 
format, or serialization is touched — this is a purely in-memory, intra-process 
concurrency primitive scoped to one `SourceFlowLifeCycle` instance's lifetime 
(a fresh `SourceCheckpointLockHandoff` per instance, no cross-restart state to 
migrate). The only observable side effect is `sourceIdleNs` metric: the handoff 
wait is now counted as idle time, whereas the previous `Thread.sleep(0L)` 
contributed nothing to that metric. This makes the metric more accurate (the 
reader genuinely is idle while ceding to an injector) but is a small, 
backward-compatible semantic refinement worth calling out explicitly since some 
existing dashboards/alerts calibrated to the old idle-time baseline could see a 
small upward shift during checkpoint injection windows — not a functional 
regression, no action required, but worth an explicit mention here for anyone 
auditing metric-based SLOs.
   
   ## 1.3 Performance / Side-Effect Analysis
   
   - Fast path (no pending injector): `awaitInjectors()` does one 
`AtomicInteger.get()` (a volatile read) and returns — negligible overhead added 
to every poll cycle, no allocation.
   - Slow path (injector pending): the reader spin-waits via `Thread.sleep(1L)` 
in a loop until the injector withdraws; bounded by how long the injector needs 
to hold the lock (documented as "normally finishes within a few milliseconds"), 
so worst-case added reader latency per checkpoint is small and self-limiting — 
it is not busy-spinning (no CPU burn between the 1ms sleeps).
   - No new locks, no new blocking I/O, no new retry logic, no new resource 
that needs releasing — the change is additive coordination state 
(`AtomicInteger`) around locking that already existed.
   - Concurrency safety: reasoned through above (finally-guarded release, 
atomic counter, no deadlock at the call sites as they exist today). I did not 
find a scenario where `pendingInjectors` can leak positive (go permanently 
non-zero without a corresponding live injector), since the only 
increment/decrement pair is co-located in `triggerBarrier` with the decrement 
in `finally`.
   
   ## 1.4 Error Handling and Logging
   
   Not applicable to the new class — 
`injectorArriving`/`injectorFinished`/`hasPendingInjector` cannot throw; 
`awaitInjectors()` correctly declares and propagates `InterruptedException` so 
that task-cancellation interrupts surface to the caller instead of being 
swallowed (verified by the new `waitingReaderHonoursInterrupt` test). No new 
logging added; none needed for this scope.
   
   **Issue 1: Test assertion window in `injectorAcquiresLockWithinOnePollCycle` 
is not fully immune to scheduling jitter, despite the comment's claim**
   - **Location**: 
`seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/task/flow/SourceCheckpointLockHandoffTest.java:106-108`
 (the `pollsBeforeInjector.set(completedPolls.get()); 
handoff.injectorArriving();` pair inside the injector thread body)
   - **Problem description**: The test comment states "The injector's 
acquisition is measured in reader poll cycles, so the assertion does not depend 
on wall-clock scheduling noise." That's the intent, but the two statements 
sampling `pollsBeforeInjector` and then calling `injectorArriving()` are not 
executed atomically with respect to the reader thread. If the injector thread 
is descheduled (GC safepoint, CI host contention) for longer than one 
`POLL_HOLD_MS` (100ms) between those two lines, the reader can complete an 
extra poll cycle in that gap, and `pollsWaited = pollsWhenInjectorAcquired - 
pollsBeforeInjector` (asserted `<= 1` at line ~139) would be inflated to 2+, 
failing the test even though the production handoff logic behaved correctly.
   - **Potential risk**: A narrow, environment-dependent flaky-test risk on 
loaded/shared CI runners — not a masked production bug (the assertion failure 
would be a false negative on the test's own timing tolerance, not evidence of 
the handoff misbehaving). Given the two statements are adjacent with no I/O or 
lock contention between them, the probability is low but not zero, especially 
under GC pressure with many parallel forked test JVMs.
   - **Best improvement**: Either (a) widen the tolerance slightly (e.g. 
`pollsWaited <= 2`) with a comment acknowledging scheduling slack, or (b) 
increase `POLL_HOLD_MS` to reduce the relative likelihood of a stray scheduling 
gap exceeding it, or (c) have `SourceCheckpointLockHandoff` (test-only 
overload, or a small test seam) return the poll-count sample atomically with 
the announcement so the measurement genuinely has no gap. I'd lean toward (b) 
as the lowest-risk change.
   - **Severity**: Medium (per Section 5.10.2 test-stability classification: 
"Risk present" — could cause intermittent failure, does not block merge, but 
should be tightened)
   - **Raised by another reviewer**: No
   
   No other issues, blocking or non-blocking, found in the production code.
   
   # 2. Code Quality Assessment
   
   ## 2.1 Coding Standards
   
   `SourceCheckpointLockHandoff` has a thorough class-level Javadoc covering 
the mechanism, the deadlock-safety argument, and the `finally`-release 
contract; every public/package method has its own Javadoc stating 
pre/postconditions (notably `awaitInjectors()`'s "the caller must not hold the 
checkpoint lock" precondition, which is the one precondition that actually 
matters for correctness). The new field/constant (`WAIT_SLICE_MS`, 
`pendingInjectors`) both have Javadoc explaining *why* their values/semantics 
are what they are, not just restating the declaration. This meets the project's 
documentation bar for concurrency-relevant fields well.
   
   ## 2.2 Test Coverage and Test Stability
   
   `SourceCheckpointLockHandoffTest.java` covers: no-op fast path, independent 
multi-injector counting, the core "bounded by one poll cycle" guarantee under 
real concurrent threads, resume-on-finish, and interrupt propagation. The class 
Javadoc explicitly ties the test back to the `BackpressureSlowSinkIT` 
regression it guards against, which is good practice.
   
   **Test-stability conclusion (Section 5.10.2)**:
   - `awaitReturnsImmediatelyWhenNoInjectorIsPending`, 
`pendingStateTracksEveryInjectorIndependently`: pure, single-threaded, 
deterministic. Stable.
   - `waitingReaderResumesWhenInjectorFinishes`: two threads, but the negative 
assertion window (reader must *not* resume within `5*WAIT_SLICE_MS+50ms`) is 
safe because nothing decrements the counter until the main thread explicitly 
does so later — no race, deterministic. Stable.
   - `waitingReaderHonoursInterrupt`: single-thread interrupt semantics, 
deterministic. Stable.
   - `injectorAcquiresLockWithinOnePollCycle`: real two-thread concurrency test 
bounded by a 10-second `JOIN_TIMEOUT_MS` safety net (so a failure surfaces as 
an assertion failure, not a CI hang) — but has the scheduling-jitter exposure 
described in **Issue 1** above.
   
   **Overall stability rating: Risk present.** The one identified pattern 
(Issue 1) is a genuine, if low-probability, source of intermittent failure 
specific to this new test; it does not indicate a flaw in the production fix 
itself, and the test correctly bounds its own worst case with a join timeout 
rather than risking a hang. Recommend tightening per Issue 1's suggestion 
before or shortly after merge.
   
   ## 2.3 Documentation Updates
   
   Not applicable — no user-visible/config-visible behavior change, so no 
`docs/en`/`docs/zh` update is required. (I did consider whether the 
metric-semantics note in 1.2 warrants a docs update; given it's a refinement of 
an existing internal metric's precision rather than a new/renamed metric or a 
config change, I don't think it clears the bar for `docs/en`/`docs/zh`, but I'm 
flagging my reasoning here in case a reviewer disagrees.)
   
   # 3. Architectural Soundness
   
   ## 3.1 Elegance of the Solution
   
   **Precise fix**, with the scope caveat documented in 1.1: it targets the 
actual mechanism of the starvation (unfair monitor re-acquisition across poll 
boundaries) without changing the public `Collector#getCheckpointLock()` 
contract, and without resorting to a blanket sleep/backoff. It is not a 
long-term solution to "bound lock hold time within one very large poll" (that 
would require either finer-grained locking inside long-running connector loops 
or a fundamentally different backpressure signal), but it correctly and 
completely solves the problem it set out to solve — the cross-poll starvation.
   
   ## 3.2 Maintainability
   
   The new class is small, single-purpose, and its Javadoc explicitly documents 
the deadlock-safety invariant future maintainers need to preserve ("the caller 
must not hold the checkpoint lock" / "release in `finally`"). The two call 
sites are minimal and easy to audit against that invariant.
   
   ## 3.3 Extensibility
   
   If a future connector pattern needs finer-grained yielding (e.g. mid-batch, 
for the JDBC-style single-lock-per-split case discussed in 1.1), 
`SourceCheckpointLockHandoff` could be extended with a second call site inside 
such a loop without changing its public contract — the design doesn't foreclose 
that.
   
   ## 3.4 Historical-Version Compatibility
   
   No checkpoint/savepoint state, API, or config surface is touched; nothing to 
migrate across upgrades.
   
   # 4. Issue Summary
   
   | # | Issue | Location | Severity |
   |---|-------|----------|----------|
   | 1 | Scheduling-jitter exposure in 
`injectorAcquiresLockWithinOnePollCycle`'s sampling window (test-stability 
"Risk present") | `SourceCheckpointLockHandoffTest.java:106-108` | Medium |
   
   # 5. Merge Recommendation
   
   ### Conclusion: Ready to merge after fixes
   
   1. **Blockers — must be fixed**: None in the production code. I'm holding 
this to "ready after fixes" rather than a bare LGTM only because of Issue 1's 
test-stability classification (Section 5.10.2 requires a `Risk present` finding 
to be listed as non-blocking-but-should-be-addressed, and I'd rather tighten a 
new test's timing assumption before it has a chance to flake on a loaded CI 
runner than fix it reactively later).
   2. **Recommended fixes — non-blocking**: Issue 1 — widen the tolerance or 
increase `POLL_HOLD_MS` in `SourceCheckpointLockHandoffTest`.
   
   **Overall assessment**: The production fix is correct and precisely scoped 
to the race it targets (I traced both call sites against `awaitInjectors()`'s 
no-lock-held precondition and found no deadlock risk, and confirmed the 
`finally`-guarded release makes starvation-of-the-reader impossible even on a 
failed injection). I want to be explicit about the one scope limitation I 
found: for connectors that hold the checkpoint lock for an entire large 
`pollNext()` batch (JDBC being a concrete example I verified in this codebase), 
this fix bounds cross-poll starvation but does not shrink the lock-hold time of 
one very large poll — that's an accurate, not overstated, claim once you read 
"bounded by a single poll" precisely, but I'd have appreciated it being spelled 
out in the class Javadoc, and I'm noting it here for anyone using this PR as a 
reference for checkpoint-latency guarantees going forward. Net: correct, 
well-tested, compatible fix; one non-blocking test-hardening item (Issue 1
 ) I'd like addressed before or shortly after merge.
   


-- 
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