DanielLeens commented on PR #11809:
URL: https://github.com/apache/seatunnel/pull/11809#issuecomment-5340851852
What changed since my last review: nothing in the code. I re-fetched the
current head (`fbe4e66d2df`) and diffed it against `dev`'s merge-base — the
production diff is byte-for-byte identical to what I reviewed on 2026-08-17
(same 3 files, +253/-6). The only new commit is another empty `[Chore][Zeta]
Retrigger PR #11809 build` commit (`git show --stat` on it lists no files). I
re-read `JobHistoryService.java`, `CoordinatorService.java`, and the new test
file end-to-end at the current head rather than trusting that "diff unchanged"
claim as a shortcut, since this round specifically asks me to re-confirm the
"3, not 1" listener-completeness claim from issue #11807 and check for a
removal-before-drain race.
# Re-verification against issue #11807's "3 listeners, not 1" claim
Confirmed directly at the current head, not inferred from the PR description:
- `JobHistoryService` registers exactly three cluster-wide
`EntryExpiredListener`s in its constructor, one each on `finishedJobStateImap`,
`finishedJobMetricsImap`, `finishedJobDAGInfoImap`
(`JobHistoryService.java:154-162`), and captures all three returned `UUID`s
into `final` fields (`:117-134`).
- `close()` (`:184-195`) calls `removeEntryListenerQuietly()` three times,
once per captured id, covering all three maps — not a subset.
- `removeEntryListenerQuietly()` (`:450-457`) wraps each
`imap.removeEntryListener(id)` in try/catch, logging at `warning` on failure,
never throwing — so a transient failure removing one listener cannot prevent
the other two from being attempted, and cannot propagate out of `close()`.
- A repo-wide sweep for `addEntryListener`/`addLocalEntryListener` confirms
these are the only three registrations tied to `JobHistoryService`'s lifecycle;
the two other engine-side `addEntryListener` call sites (in the
checkpoint/metrics state stores) already have their own correct capture/remove
symmetry and are unrelated to this fix.
So: all three leaking listeners identified in #11807 are captured and
removed by this PR, not just one. No partial-removal or double-removal issue —
`IMap.removeEntryListener` on an already-removed id returns `false` rather than
throwing, so calling `close()` twice (which can't currently happen, since it's
gated by the CAS below, but is worth noting as safe-by-construction) would be a
no-op the second time.
# Re-verification against the "removal-before-drain race" /
abnormal-termination question
- The single production teardown call site is
`CoordinatorService.clearCoordinatorService()`
(`CoordinatorService.java:1211-1213`, `if (jobHistoryService != null) {
jobHistoryService.close(); }`), which is itself `synchronized` and gated by
`coordinatorServiceCleared.compareAndSet(false, true)` (`:1174`).
`initCoordinatorService()` resets that flag as its first statement (`:516`), so
every active→inactive cycle re-arms the guard and `close()` runs exactly once
per cycle — not zero times, not twice.
- `clearCoordinatorService()` is reached on **all three** relevant paths,
not only normal completion: role loss (`checkNewActiveMaster` Case A),
`initCoordinatorService()` throwing during startup (Case C, caught at
`:1152-1160` — the guard being freshly reset at `:516` means a failed partial
init still gets cleaned up), and full node `shutdown()` (`:1966-1977`). So
listener deregistration fires on abnormal termination and failure paths too,
not just the happy path.
- Hazelcast does not quiesce in-flight `entryExpired()` callbacks on
`removeEntryListener()`, so an event that started dispatching just before
`close()` can still complete concurrently. That's benign here:
`incrementFinishedJobCleanupTotal` uses a `ConcurrentHashMap<String,
AtomicLong>` (`:115` / `:440-442`), and the `CleanLogOperation` fan-out is
already wrapped in its own try/catch, so a concurrent in-flight callback can't
corrupt state or throw out of the listener thread.
- The one genuine timing gap — which I flagged in my 08-16 round and is
still accurate and still undocumented — is the opposite direction: a window
with **zero** registered listeners between the outgoing master's `close()` and
the incoming master's fresh registrations in `initCoordinatorService()`. An
expiry landing in that exact window is silently dropped (no reconciliation
sweep exists). This is Low severity (bounded by the up-to-20s
`awaitTermination` that runs before `close()`, and Hazelcast TTL expiry being
lazy), not a correctness blocker, but the `close()` Javadoc still doesn't
mention it.
# Current issue list (unchanged from my 2026-08-17 review — carried forward,
not re-numbered)
| # | Issue | Location | Severity | Status |
|---|-------|----------|----------|--------|
| 1 | Constructor assigns the three listener ids to `final` fields inline; a
partial-construction failure (2nd/3rd `addEntryListener` throwing) leaks the
already-registered listener(s) unremovably | `JobHistoryService.java`
constructor (`:154-162`) | Low | Unresolved, 3rd round |
| 2 | Zero-listener handover window (described above) is undocumented in the
`close()` Javadoc | `CoordinatorService.java:1211-1213`,
`JobHistoryService.java:184-195` | Low | Unresolved, 3rd round |
| 3 | `testClearCoordinatorServiceDeregistersJobHistoryListeners` calls
`coordinatorService.getJobHistoryService()`/`seaTunnelServer.getCoordinatorService()`
with no `await()` wrapper. I re-read `SeaTunnelServer.getCoordinatorService()`
at the current head (`SeaTunnelServer.java:286-317`): if the node isn't master
yet it throws immediately with zero retries; if it is master but the
coordinator hasn't finished initializing, the internal retry is capped at
3×500ms=1.5s before throwing. The sibling test file already establishes the
correct idiom for this exact call (`await().atMost(60,
TimeUnit.SECONDS).untilAsserted(...)`, used 8+ times in
`CoordinatorServiceTest.java`), and the new test doesn't follow it |
`JobHistoryServiceListenerCleanupTest.java:104-119` | Medium | Unresolved, 3rd
round |
| 4 | The isolated Hazelcast instance in the new test is created via
`SeaTunnelServerStarter.createHazelcastInstance(clusterName)` with only the
cluster name overridden — default network/port config. The sibling test file
has an existing helper,
`createHazelcastInstanceWithJoinPortTryCount(clusterName, 100)`, specifically
to avoid join/port contention when multiple node instances run in the same CI
process; this test (which extends `AbstractSeaTunnelServerTest` and therefore
already has a base-class node running) doesn't use it |
`JobHistoryServiceListenerCleanupTest.java:105-108` | Medium | Unresolved, 3rd
round |
| 5 | Neither test asserts a cross-instance isolation case (closing one
`JobHistoryService` must not remove a different, still-live instance's
registrations) or an absolute listener-count on the maps — an over-removing
`close()` would still pass this suite |
`JobHistoryServiceListenerCleanupTest.java:60-94` | Low | Unresolved, 3rd round
|
| 6 | Non-`synchronized` `initCoordinatorService()` can theoretically
register fresh listeners after the shutdown thread has already run `close()`.
Pre-existing structural asymmetry, not introduced by this PR, out of scope here
| `CoordinatorService.java:515` vs `:1172` | Low | Unresolved, out of scope |
No new issues found in this pass. The production logic conclusion is
unchanged and, on this fourth independent read, still holds: the fix is
complete with respect to all three listeners named in #11807, is reached on the
normal-role-loss, init-failure, and shutdown paths (not only the "happy path"),
and does not introduce a double-removal or a removal-before-drain correctness
bug.
# CI status (checked live at review time)
The apache-side `Build` check still shows `pending` — per the standing
project fact, that's only a pointer, so I went to the actual fork run
(`DanielLeens/seatunnel`, run for head `fbe4e66d2df`). All 30 jobs have
completed except the run's own top-level status, which is still reporting
`queued` — that looks like an API-side lag rather than real pending work, since
I could not find any job still in `queued`/`in_progress` state.
Results:
- `unit-test (8, ubuntu-latest)`, `unit-test (8, windows-latest)`: pass.
- `unit-test (11, ubuntu-latest)`: **cancelled** (not evaluable on its own;
needs a rerun).
- `unit-test (11, windows-latest)`: **failure** —
`TaskExecutionServiceTest.before` → `IllegalStateException: Node failed to
start!` after 311.37s.
- All `all-connectors-it-*` and `jdbc-connectors-it-part-*` legs that ran:
pass.
I pulled the raw log for the failing job and specifically checked whether
this PR's new test is implicated, since Issue 4 above is exactly about
port/network isolation. It isn't: `JobHistoryServiceListenerCleanupTest` ran
earlier in the same job and passed cleanly (2/2, 2.972s), its Hazelcast
instances logged a clean `SHUTDOWN` before the next class started, and the
failing class (`TaskExecutionServiceTest`) is unrelated code this PR doesn't
touch, in a different package. This matches a `windows-latest` engine-server
"Node failed to start!" flake pattern that predates this PR and has recurred on
unrelated PRs hitting a different test class each time — it reads as
environmental, not caused by this diff.
Given project policy, an unverified/failing CI still blocks merge regardless
of suspected cause. Recommended next step: rerun just the two JDK-11
`unit-test` jobs (job-level rerun, not a full workflow rerun) rather than
treating this as a code problem.
# Merge recommendation
### Conclusion: Not ready to merge yet — source logic is correct and
complete, but CI needs a clean rerun and the two Medium test-stability issues
should be fixed now rather than carried to a 4th round.
1. **Blockers — must be fixed before merge**
- CI: `unit-test (11, windows-latest)` failed (environmental flake,
unrelated to this diff on the evidence above) and `unit-test (11,
ubuntu-latest)` was cancelled. Both need a clean pass before merge per project
policy.
- Issue 3 (Medium): add the `await().atMost(...).untilAsserted(...)`
guard around the coordinator-service acquisition in
`testClearCoordinatorServiceDeregistersJobHistoryListeners`, matching the
existing idiom in `CoordinatorServiceTest.java`. Three review rounds in and
this is still a 3-line fix using a pattern that already exists in the same test
package.
- Issue 4 (Medium): build the isolated Hazelcast instance from
`createHazelcastInstanceWithJoinPortTryCount(...)` instead of the
default-network `createHazelcastInstance(clusterName)`, matching the same
sibling file.
2. **Recommended fixes — non-blocking**
- Issue 1 (Low): make the three constructor registrations all-or-nothing,
or switch to the `volatile UUID` + null-guard idiom the sibling state-store
classes already use.
- Issue 2 (Low): document the zero-listener handover window in the
`close()` Javadoc.
- Issue 5 (Low): add a cross-instance isolation assertion so an
over-removing `close()` would be caught.
- Issue 6 (Low): out of scope for this PR, noted only.
**Overall assessment.** On a fourth independent pass through the production
code, the fix remains correct, minimal, and complete: all three listeners named
in issue #11807 are captured and removed, the CAS guard re-arms correctly so
cleanup runs exactly once per transition, and the single teardown call site is
reached on role-loss, init-failure, and shutdown, not only normal completion. I
found no removal-before-drain race and no double-removal risk —
`IMap.removeEntryListener` is idempotent and the try/catch in
`removeEntryListenerQuietly` prevents one failed removal from blocking the
other two. What's holding this back is unchanged from my last round: CI isn't
clean yet (one apparently-environmental Windows flake plus one cancelled job
needing a rerun), and the two Medium test-stability findings — both cheap, both
following patterns that already exist verbatim in the same test file's sibling
class — still haven't been picked up across three rounds of review. I'd rather
fix
Issues 3 and 4 now and get a clean CI run than keep re-confirming the same
two open items on a fourth pass.
(Note: since I authored this PR myself, this is posted as a plain comment
rather than a formal approve/request-changes review — GitHub doesn't allow
self-approval, and I want the same level of independent scrutiny applied here
as to anyone else's PR.)
--
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]