DanielLeens commented on PR #12316:
URL: https://github.com/apache/seatunnel/pull/12316#issuecomment-5846335176
Thanks for the thorough pass — I went back through the current head against
each point rather than taking my own prior conclusions for granted, and I want
to respond to all five with where I actually landed.
**Issues 1 & 2 (timer flush signal path not covered by the handoff)**:
you're right, and I should have caught this myself. I re-traced it:
`onTimerTick` (`SourceFlowLifeCycle.java:204-214`) calls
`collector.sendFlushSignal(...)`, which goes through `sendRecordToNext` and
`synchronized (checkpointLock)` (`SeaTunnelSourceCollector.java:410`) — the
identical monitor `triggerBarrier` registers against, but `onTimerTick` never
calls `injectorArriving()`/`injectorFinished()`. So the reader's
`awaitInjectors()` has no visibility into a pending flush tick, and that tick
is exposed to exactly the unfair-reentry race this PR eliminates for barriers.
It's pre-existing behavior, not a regression from this diff, but since
`timerFlushWorker` is a shared, node-wide pool, a tick losing that race for
several multi-second polls does pin a shared thread longer than necessary, and
it's a gap in what this PR set out to fix. I'll wrap `onTimerTick`'s
`sendFlushSignal` call in the same `injectorArrivi
ng()` / `try { ... } finally { injectorFinished(); }` pattern in the next
revision — likely via a small shared helper on `SourceCheckpointLockHandoff` so
both call sites can't drift apart again — rather than just documenting the gap
away.
**Issue 3 (sampling order in `injectorAcquiresLockWithinOnePollCycle`)**:
confirmed, and your fix is better than what I originally proposed. I had
flagged the same test in my 09-15 review and suggested widening the tolerance
or raising `POLL_HOLD_MS`, which only shrinks the window. Working through it
again with your suggested swap — announce via `injectorArriving()` first, then
sample `completedPolls` — actually removes the race: once the injector has
registered, the reader's very next `awaitInjectors()` call is guaranteed to see
it and park before it can re-acquire the lock for an unannounced extra poll, so
there's no gap left to sample around. I'll make that swap and correct the
comment above it.
**Issue 5 (javadoc wording)**: agreed, "an injector never waits for the
reader" is only true in the deadlock-avoidance sense and reads like a latency
claim on its own. I'll tighten that sentence and add an explicit note on which
`checkpointLock` contenders are (and, until Issue 1 lands, aren't) registered
with the handoff, so this doesn't have to be inferred from the call sites.
**Issue 4 (sleep-poll vs. a signalled wait)**: this one I'd like to hold as
a follow-up rather than fold into this PR. The design point is fair — a
`wait`/`notifyAll` (or `LockSupport.park`/`unpark`) handoff would resume the
reader immediately instead of on the next 1ms slice — but it trades a very
small, well-understood latency cost for a different correctness surface to
reason about (missed-notify ordering between `injectorFinished()` and a reader
that hasn't started waiting yet). Given the documented cost is already bounded
and this PR is meant to stay a minimal, narrowly-scoped fix, I'd rather land
the flush-path coverage and the test fix first and revisit the wait mechanism
separately if there's ever measured evidence the 1ms granularity matters in
practice. Happy to hear if you disagree on the priority.
I'll push a commit covering Issues 1, 2, 3 and 5 next; once that's up I'd
like you to take another look specifically at the flush-path wiring, since
that's the part that's genuinely new rather than a restatement of what I'd
already checked.
--
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]