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]