DanielLeens commented on PR #11809:
URL: https://github.com/apache/seatunnel/pull/11809#issuecomment-5315376167
# What Problem Does This PR Solve?
`JobHistoryService` registers three cluster-wide Hazelcast
`EntryExpiredListener`s in its constructor but never deregistered them, so
every active-master transition leaked a full set of listeners. This PR captures
the three registration ids and removes them from
`CoordinatorService.clearCoordinatorService()`.
Concretely what accumulates: `CoordinatorService.initCoordinatorService()`
constructs a brand-new `JobHistoryService` every time this node becomes the
active master. That constructor registers three entry listeners on three
cluster-wide IMaps (finished-job state, finished-job metrics, finished-job
vertex info). Before this PR the return value of `addEntryListener(...)` was
discarded, and no code path ever called `removeEntryListener`. Because both
listener classes are non-static inner classes, each abandoned registration
keeps its enclosing `JobHistoryService` — and transitively the node engine, the
IMaps, and the running/pending job maps — reachable from the IMap listener
registry, and one of the listeners is the sole producer of a cleanup RPC in the
codebase, so leaked registrations also duplicate that RPC fan-out on every real
expiration event.
Fix approach: store the three UUIDs returned by `addEntryListener` in new
final fields, add a `close()` that deregisters exactly those three via a
fail-soft helper, and invoke it from
`CoordinatorService.clearCoordinatorService()`. `jobHistoryService` is
deliberately not nulled, so read paths keep working after the role switch.
# What Changed Since My Last Review
Nothing in the code. I re-fetched `upstream/dev`, diffed the current head
against the merge-base, and confirmed the production diff is byte-for-byte
identical to what I reviewed on 2026-08-16 (the same 3 files changed). The only
commit added since my last comment is an empty "retrigger build" commit — `git
show --stat` on it lists no files. None of my six previously-raised issues (2
Medium, 4 Low) have been addressed in source. I re-verified this directly
against the current head files rather than trusting the diff being empty as a
shortcut.
# 1. Code Change Review
## 1.1 Core Logic Analysis
I re-read `JobHistoryService.java` and `CoordinatorService.java` in full at
the current head (not the diff) and re-independently re-verified the four
load-bearing claims from my prior round, because a re-review that just restates
prior conclusions without re-checking is worthless:
- **All three listeners captured and removed** — confirmed line-for-line
against current head.
- **Single production construction site** — remains the only `new
JobHistoryService(` in non-test code.
- **CAS guard re-arms correctly** — the guard flag is reset as the first
statement of `initCoordinatorService()`, and the CAS gate in
`clearCoordinatorService()` is unchanged — every active/inactive cycle still
performs exactly one cleanup.
- **`close()` call site unchanged** — still inside the `synchronized`
`clearCoordinatorService()`, still reached from role-loss, init-failure-retry,
and shutdown.
**Runtime path (unchanged from prior round, re-verified against current line
numbers):**
```
Case A leave active master
-> clearCoordinatorService() [synchronized]
CAS(false->true) guard
awaitTermination(20s)
jobHistoryService.close() <-- THE FIX
-> removeEntryListenerQuietly x3
resourceManager.close() ; eventProcessor.close()
Case B become active master
-> initCoordinatorService()
coordinatorServiceCleared.set(false) <-- re-arms
jobHistoryService = new JobHistoryService(...) -> 3 fresh ids
Case C init throws -> clearCoordinatorService(), guard is false, close()
runs
Case D shutdown() -> clearCoordinatorService() -> same close() path
```
**Normal job lifecycle does not reach this code, and that is correct.**
Submit -> run -> finish never touches `clearCoordinatorService()`; the removal
fires only on master-role loss, init failure, and shutdown — not on job
completion despite the PR title's phrasing. This is the correct scope:
cancelling a job must not tear down the master's history listeners, and I
confirmed no other code path calls `close()`.
## 1.2 Compatibility Impact
Fully backward compatible — unchanged conclusion, re-verified. No `Option`
added/removed/renamed. `close()` and the new id-accessor are purely additive,
the latter package-private. `JobHistoryService` is not `Serializable`; the
fields that are actually serialized are untouched. No checkpoint/savepoint/IMap
payload format is affected. Rolling upgrade and downgrade are both safe — a
downgrade simply restores the pre-fix leak, no data-format break.
## 1.3 Performance / Side-Effect Analysis
This is the point of the fix, so I re-verified rather than restated:
- **Net positive, confirmed:** listener count is bounded at 3 regardless of
failover count; duplicate cleanup-RPC fan-out and duplicate cleanup-counter
increments stop after the first master switch post-fix.
- **New blocking work inside `synchronized clearCoordinatorService()`:**
unchanged from prior round. Three additional synchronous
`IMap.removeEntryListener` calls now sit inside the same synchronized method
that already blocks on a 20s `awaitTermination` and on
`resourceManager.close()`. Same risk class as those neighbors, not a new class
of risk. Not a blocker.
- **Handover window with zero registered listeners (raised in my prior
round, still not addressed in code or Javadoc):** because
`addEntryListener(listener, true)` is cluster-wide, there is a window between
the outgoing master's `close()` and the incoming master's
`initCoordinatorService()` where no listener is registered on the finished-job
IMaps at all. An expiry landing in that window is dropped permanently — no
reconciliation sweep exists to recover it. I re-checked whether anything since
has added a periodic reconciliation pass — it has not. Risk remains Low
(bounded by the up-to-20s `awaitTermination` before `close()` runs, and
Hazelcast TTL expiry being lazy), and I still do not consider it a merge
blocker, but it remains completely undocumented in the `close()` Javadoc, which
a future maintainer reading only that Javadoc would not learn about.
## 1.4 Error Handling and Logging
Unchanged and still correct: the removal helper catches `Exception`, logs at
`warning` with the map name, never propagates. The null-guard on
`jobHistoryService` covers the first-ever clear before any init. No sensitive
data logged.
No new issues found in this pass beyond the two carried over from the prior
round:
**Issue 1** *(carried over, unresolved)*
- **Location:** `JobHistoryService` constructor; contrast the safer pattern
used by a sibling state-store class in this module.
- **Problem:** The three ids are `final` fields assigned inline. If the
second or third `addEntryListener` call throws, the already-registered
listener(s) become permanently unremovable — the exact leak this PR fixes, in
one uncovered corner. Because `initCoordinatorService()` failure is retried
every ~100ms by the active-master check, a persistent partial-construction
failure would leak on every tick.
- **Risk:** Narrow and pre-existing in shape (this exact hole existed before
the PR too, just without any cleanup path at all), not introduced by this PR,
but not closed by it either.
- **Severity:** Low. Not required for merge.
**Issue 2** *(carried over, unresolved)*
- **Location:** `CoordinatorService.clearCoordinatorService()` interacting
with `JobHistoryService`'s registration/close methods.
- **Problem:** Zero-listener handover window described in 1.3, undocumented
in the `close()` Javadoc.
- **Risk:** Low probability, orphaned log files on member nodes if an expiry
lands exactly in the window.
- **Best improvement:** No code change required for merge; add a sentence to
the `close()` Javadoc so the window is a documented, intentional trade-off
rather than something the next maintainer has to rediscover from first
principles.
- **Severity:** Low.
# 2. Code Quality Assessment
## 2.1 Coding Standards
Unchanged and re-verified: ASF license header present and correct on the new
test file. All three new fields and all three new methods carry multi-line
Javadoc explaining purpose, not just type. No wildcard imports, no
`System.out.println`, no single-line Javadoc.
## 2.2 Test Coverage and Test Stability — does the test actually detect a
leak, and is it CI-safe?
**Yes, it genuinely detects a leak, not just "no exception thrown."** I
re-read the test file at the current head end-to-end. Both tests assert on
`IMap.removeEntryListener(UUID)`'s boolean return, which is a real observation
of the listener registry: `true` means the registration is live, `false` means
it is gone. One test establishes a positive control first — proving the exposed
ids are real live registrations — before running three create/close cycles and
asserting `false` after each, which directly emulates repeated
master-transition churn. The other test exercises the real production entry
point (`clearCoordinatorService()`) rather than calling `close()` directly,
which is the correct level to test at.
**Both CI-stability defects I flagged in my prior round are still present,
verbatim, in the current head — I re-read the file line-by-line to confirm, not
just trusted my earlier finding:**
- **Issue 3 (Medium, carried over, unresolved).** The test acquiring the
coordinator service has no `await()` wrapper around it. I re-read the accessor
directly at the current head: if the node isn't master yet, it throws
immediately with zero retries; if it is master, the internal retry loop is
bounded to 1.5s before throwing. The module's own sibling test file establishes
the answer to exactly this problem — an `await().atMost(60,
TimeUnit.SECONDS).untilAsserted(...)` idiom used repeatedly throughout that
file. The new test does not follow this established idiom. On a loaded CI
runner, single-node master election plus `initCoordinatorService()` completing
inside 1.5s is not guaranteed, and this project already has a documented
flaky-engine-server-test burden.
- **Issue 4 (Medium, carried over, unresolved).** The isolated Hazelcast
instance in the new test is built with only the cluster name overridden,
inheriting default network/port settings. I re-confirmed the sibling test file
has a private helper used at 8+ call sites specifically to avoid port-join
contention between concurrently-running node instances in the same CI job. This
test class also stands up a second node alongside the base-class node without
that same protection.
Both are exactly as I described them in my prior round; no attempt has been
made to address either, and their fix cost is genuinely small (three lines
each, using patterns that already exist verbatim elsewhere in the same test
package).
**Coverage gaps also carried forward, unaddressed:** no absolute
listener-count assertion, no cross-instance isolation assertion (an
over-removing `close()` that killed a concurrent instance's listeners would
still pass this suite), no coverage of the shutdown path or init-failure path.
**Stability rating: acceptable but should be hardened before merge —
unchanged from my prior conclusion.** The fundamentals (no `Thread.sleep`, no
floating-point tolerance, no inter-test ordering dependency, `finally`-guarded
teardown) are genuinely good. But Issues 3 and 4 are real,
established-pattern-diverging test-stability risks being introduced into a
module with a known flaky-test history, and they remain unfixed in this round.
## 2.3 Documentation Updates
None required, none made — correct, unchanged.
# 3. Architectural Soundness
## 3.1 Elegance of the Solution
Unchanged: a precise root-cause fix restoring the missing
register/deregister symmetry at the one construction site and one teardown
site, with no new abstraction or config. Almost entirely additive.
## 3.2 Maintainability
Unchanged: `close()` sits alongside the existing resource-teardown idiom in
`clearCoordinatorService()` and reads as the established pattern for
coordinator-owned-resource teardown.
## 3.3 Extensibility
Unchanged, non-blocking: the three `(IMap, UUID, storeName)` triples are
three parallel fields rather than one collection. Fine at the current fixed
count of three; worth generalizing only if a fourth listener is added later.
## 3.4 Historical-Version Compatibility
Unchanged: registration ids are ephemeral in-memory state, never persisted,
never part of a checkpoint/savepoint payload. Rolling upgrade and downgrade are
both safe.
# 4. Issue Summary
| # | Issue | Location | Severity | Status |
|---|-------|----------|----------|--------|
| 1 | Constructor assigns three listener ids to `final` fields inline; a
partial-construction failure leaks the already-registered listener(s)
unremovably | `JobHistoryService.java` (constructor) | Low | Carried over,
unresolved |
| 2 | Zero-listener handover window is undocumented in the `close()` Javadoc
| `CoordinatorService.java`, `JobHistoryService.java` | Low | Carried over,
unresolved |
| 3 | Coordinator-service accessor in the new test has no Awaitility guard,
diverging from the established sibling-test pattern used at 8+ sites in the
same module | `JobHistoryServiceListenerCleanupTest.java` | Medium | Carried
over, unresolved |
| 4 | Isolated Hazelcast instance in the new test built from default network
config, unlike the port-contention-safe pattern used throughout the sibling
test file | `JobHistoryServiceListenerCleanupTest.java` | Medium | Carried
over, unresolved |
| 5 | No cross-instance isolation assertion and no absolute listener-count
assertion | `JobHistoryServiceListenerCleanupTest.java` | Low | Carried over,
unresolved |
| 6 | Non-synchronized `initCoordinatorService()` can theoretically register
fresh listeners after the shutdown thread already ran `close()`; pre-existing
structural asymmetry, not introduced by this PR | `CoordinatorService.java` |
Low | Carried over, out of scope for this PR |
No new issues found in this re-review beyond the six carried over from my
last round.
# 5. CI Status
Re-checked live at review time, not inferred from a stale snapshot: the
apache-side pointer showed `FAILURE`, but per the standing project fact that
the apache-side check is only a pointer, the real signal is the fork run. I
pulled the full job list and the failing job's raw log: `unit-test` passed on
all four matrix legs (the module most directly exercising this change),
`engine-v2-it` passed on both JDKs. The single failure is one connector-IT job,
and I pulled the actual Maven output for it: the failure is a Testcontainers
container-startup failure for a connector this PR does not touch (the diff is
confined to `seatunnel-engine`). This reads as an environment/infra flake, not
a regression caused by this change, but per project policy the PR still needs a
green CI run before merge; this failure has not yet been re-verified as
transient by a rerun.
# 6. Merge Recommendation
### Conclusion: Not ready to merge. Source logic remains correct; CI is red
on an apparently-unrelated infra flake, and two Medium test-stability issues
from my prior review are still unaddressed one round later.
1. **Blockers — must be fixed before merge**
- **CI is red.** The fork run's connector-IT job failed on a
Testcontainers startup issue, unrelated to this diff on its face, but per
project policy an unverified/failing CI blocks merge regardless of cause. This
needs a rerun to confirm it is transient before this PR can be considered
CI-clean.
- **Issue 3 (Medium):** wrap the coordinator-service acquisition in the
new test in the same `await().atMost(...).untilAsserted(...)` idiom already
used 8+ times in the sibling test file. This is the most likely source of a
future flaky-CI report on this exact new test, and remains a three-line fix not
yet made across two review rounds.
- **Issue 4 (Medium):** build the isolated instance in the new test from
an explicit config with an allocated port and raised join-port-try-count,
mirroring the sibling test's existing helper.
2. **Recommended fixes — non-blocking**
- Issue 1 (Low): make the three constructor registrations all-or-nothing,
or adopt the volatile-field + null-guard idiom the sibling state stores already
use.
- Issue 2 (Low): document the zero-listener handover window in the
`close()` Javadoc.
- Issue 5 (Low): add a cross-instance isolation assertion.
- Issue 6 (Low): out of scope, note only.
**Overall assessment.** The production fix itself remains correct, minimal,
and well-reasoned on this third independent pass: I re-derived the leak
mechanism, re-confirmed the CAS re-arm, re-confirmed there is exactly one
production construction site and one teardown site, and re-confirmed
compatibility is clean on every axis (no API, no config, no
serialized/checkpoint state). Nothing about the source correctness conclusion
has changed since my last round, and I found no new production-code issues in
this pass. What has not changed is more important right now: this PR sat for a
full day with CI queued, then went green everywhere that matters to the actual
bug except a plausibly-unrelated Testcontainers flake, and neither of the two
Medium test-stability findings from my last review — both three-line fixes,
both following patterns that already exist in the very same test file's sibling
class — were acted on before the retrigger. Holding my own PR to the full bar
means not lett
ing "CI will probably go green on rerun" substitute for actually fixing
findings I already wrote down and already know how to fix. I'd fix Issues 3 and
4 now, rerun the failed job, and only then move this to mergeable.
(Note: since I authored this PR myself, this review 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]