DanielLeens commented on PR #11809:
URL: https://github.com/apache/seatunnel/pull/11809#issuecomment-5378208862
This is my own PR, so per protocol I'm posting this as a plain issue comment
(GitHub blocks self-approval/self-request-changes via the review API) and
stating the merge determination explicitly in the text.
Status check before the substance: I re-fetched the PR live. The head is
still `3851e0d9f2b204fa0cb7d030d9d64321dbccd936` (unchanged since my last round
yesterday at 2026-08-21T14:13:03Z), and no one has posted a new commit,
comment, or review since then. So there is nothing new to react to — but since
this is a self-authored PR, protocol asks me to re-do a full, independent read
of the current diff rather than just restate the prior conclusion from memory.
I did that below, re-reading `CoordinatorService.java`,
`JobHistoryService.java`, and the new test file directly at the current head,
and I land in the same place as yesterday's round.
# What Problem Does This PR Solve?
- User pain point: `JobHistoryService` registers three cluster-wide
Hazelcast `EntryExpiredListener`s (on `IMAP_FINISHED_JOB_STATE`,
`IMAP_FINISHED_JOB_METRICS`, `IMAP_FINISHED_JOB_VERTEX_INFO`) in its
constructor, and `CoordinatorService.checkNewActiveMaster()` ->
`initCoordinatorService()` constructs a brand-new `JobHistoryService` every
time this node becomes the active master (`CoordinatorService.java:515-529`).
Before this PR nothing ever deregistered the previous instance's listeners, so
every active -> inactive -> active master transition (failover, split-brain
healing, rolling restart) leaked three listener registrations, kept the old
`JobHistoryService` (and everything it transitively references) reachable from
the IMap listener registry, and caused a real finished-job expiration event to
fire once per leaked instance — duplicating `CleanLogOperation` RPC fan-out and
cleanup counters.
- Fix approach: capture the `UUID` returned by each `addEntryListener(...)`
call as `final` fields, add a best-effort `JobHistoryService.close()`
(`JobHistoryService.java:184-195`) that deregisters exactly those three via a
`removeEntryListenerQuietly` helper (`:450-457`, catches and logs, never
throws), and invoke `close()` from
`CoordinatorService.clearCoordinatorService()`
(`CoordinatorService.java:1211-1213`) — the single method already used both
when a node leaves the active-master role and during full coordinator shutdown.
`jobHistoryService` is deliberately left non-null so in-flight readers keep
working after the role switch.
- One-sentence summary: closes a listener/side-effect leak on Zeta
master-role transitions by tracking and explicitly deregistering the three
finished-job expiration listeners when the coordinator service is cleared.
# 1. Code Change Review
## 1.1 Core Logic Analysis
**Before** (`JobHistoryService.java` constructor):
```java
this.finishedJobStateImap.addEntryListener(
new FinishedJobExpiredListener<>(Constant.IMAP_FINISHED_JOB_STATE),
true);
this.finishedJobMetricsImap.addEntryListener(
new
FinishedJobExpiredListener<>(Constant.IMAP_FINISHED_JOB_METRICS), true);
this.finishedJobDAGInfoImap.addEntryListener(
new JobInfoExpiredListener(Constant.IMAP_FINISHED_JOB_VERTEX_INFO),
true);
```
Registration ids discarded, no removal path anywhere.
**After** (`JobHistoryService.java:154-162`):
```java
this.finishedJobStateListenerId =
this.finishedJobStateImap.addEntryListener(
new
FinishedJobExpiredListener<>(Constant.IMAP_FINISHED_JOB_STATE), true);
this.finishedJobMetricsListenerId = ...
this.finishedJobDAGInfoListenerId = ...
```
plus `close()` and, in `CoordinatorService.clearCoordinatorService()`:
```java
if (jobHistoryService != null) {
jobHistoryService.close();
}
```
**Runtime flow (verified directly against the current head, not assumed):**
```text
Normal path (single-threaded scheduler tick, no race):
masterActiveListener (1-thread ScheduledExecutorService, 100ms tick)
-> CoordinatorService.checkNewActiveMaster()
CoordinatorService.java:1133
node loses master role
-> clearCoordinatorService()
CoordinatorService.java:1150/1157, synchronized
executorService.shutdownNow() + awaitTermination(20s)
:1189-1201
-> jobHistoryService.close() :1211-1213 (new)
-> removeEntryListenerQuietly x3
JobHistoryService.java:184-195 -> :450-457
-> resourceManager (capture-and-null) .close() :1216-1219
-> eventProcessor (capture-and-null) .close() :1222-1228
Node shutdown path (the one narrow exception, see Issue 1):
CoordinatorService.shutdown() :1966
masterActiveListener.shutdown() (non-blocking on an in-flight tick)
:1969
clearCoordinatorService() (called directly from the shutdown
caller's thread) :1977
awaitSchedulerTermination(...) (only awaited AFTER
clearCoordinatorService runs) :1978
```
**Key Findings:**
- I independently re-verified the core fix against the current code (not
just the diff): all three listener ids are captured as `final UUID` fields
(`JobHistoryService.java:122/128/134`), `close()` removes exactly those three
via `removeEntryListenerQuietly` (`:184-195`), and the helper swallows
exceptions per-map so a transient failure removing one listener cannot block
the other two (`:450-457`). This matches the leak description in issue #11807
exactly and is the correct fix for the normal, scheduler-driven transition path.
- `clearCoordinatorService()` is `public synchronized void`
(`CoordinatorService.java:1172`); `initCoordinatorService()`, which reassigns
the `jobHistoryService` field (`:529`), is **not** synchronized on the same
monitor and the field is not `volatile`. Under the routine path both are only
ever invoked from the single `masterActiveListener` scheduler thread via
`checkNewActiveMaster()`, so they can't interleave there. But `shutdown()`
(`:1966-1979`) calls `masterActiveListener.shutdown()` — which does not
interrupt or wait for a tick already in flight — and then calls
`clearCoordinatorService()` directly from the calling thread before ever
awaiting scheduler termination. If a `checkNewActiveMaster()` tick is
mid-flight on the scheduler thread at that exact moment, between
`coordinatorServiceCleared.set(false)` (`:516`, which re-arms
`clearCoordinatorService()`'s idempotency guard) and the `jobHistoryService =
new JobHistoryService(...)` assignment (`:529`), the shutdown-calle
r thread's unsynchronized read of `jobHistoryService` can observe the freshly
published instance and close **that** one instead of the outgoing generation's.
This is a real, code-verifiable data race (two threads touching an
unsynchronized, non-volatile field with no shared lock), not a theoretical
concern — I traced every step above directly in the current source rather than
trusting the earlier description of it.
- The three-line block that deregisters `jobHistoryService` does **not**
follow the capture-and-null idiom already used three lines below it for
`resourceManager` (`CoordinatorService.java:1216-1219`, `ResourceManager
manager = resourceManager; resourceManager = null;`) — that idiom exists
precisely to make this class of hazard impossible for other fields in the same
method, and the new code doesn't reuse it.
- `JobHistoryService`'s constructor still performs all three
`addEntryListener` calls unguarded (`:154-160`); if the second or third throws
mid-construction, the instance is never published and any id(s) already
obtained on the earlier map(s) are permanently orphaned. This is not a
regression (pre-fix, every construction leaked all three unconditionally), but
it's a residual gap in the same leak class the PR targets.
- The class does not declare `implements AutoCloseable` even though it now
has a `close()` matching that signature exactly (`:69`), and the test-only
accessor `getEntryListenerRegistrationIds()` (`:202`) is package-private but
not marked `@VisibleForTesting`, despite that annotation already being used
elsewhere in this module (`CoordinatorService.java`,
`CheckpointCoordinator.java`, `PhysicalVertex.java`, `JobMaster.java`).
## 0.5 Incorporating other reviewers
`@SEZ9` posted a substantive review on 2026-08-21T06:54:03Z with 8 numbered
findings. I re-verified each against the current head myself rather than
trusting the summary:
- Issue 1 (mutable-field race on
`clearCoordinatorService()`/`initCoordinatorService()`): confirmed
independently above — real, Medium. Agree with `@SEZ9`.
- Issue 2 (closed-but-non-null field could be reused by a future refactor):
I checked — `initCoordinatorService()` today always does `jobHistoryService =
new JobHistoryService(...)` unconditionally, with no `if (jobHistoryService ==
null)` guard anywhere, so this is not an active bug today, only a
forward-looking fragility. I rate it Low rather than `@SEZ9`'s Medium; `@SEZ9`,
if you disagree I'm happy to hear the concrete path where it bites today.
- Issue 3 (partial-construction orphan): confirmed the code path is real,
but since it's strictly better than the pre-fix baseline (which leaked all
three unconditionally, every time) and requires a rare exception
mid-construction, I rate it Low rather than Medium — happy to be overruled if
there's a concrete production trigger I'm missing.
- Issues 4-7 (test try/finally hygiene, swallowed removal-failure
diagnostics, `AutoCloseable`, `@VisibleForTesting`): confirmed all four
directly against the current code (test file line 77 `liveService` has no
try/finally; `removeEntryListenerQuietly` discards the boolean from
`removeEntryListener`; the class has no `implements AutoCloseable`; the
accessor has no `@VisibleForTesting`). Agree, Low severity on all four, +1 to
`@SEZ9`.
- Issue 8 (raw-type test class): checked the surrounding test suite's actual
convention — the closest siblings in purpose
(`CoordinatorServiceJobCleanupTest`,
`CoordinatorServiceWithCancelPendingJobTest`,
`CoordinatorServicePipelineCleanupTest`, all `CoordinatorService`-lifecycle
tests like the new one) also extend `AbstractSeaTunnelServerTest` raw, with no
type argument. I would not block on this — it matches its closest siblings,
even if other, unrelated tests in the module use the generic form.
None of `@SEZ9`'s findings have been addressed by a new commit — the head is
unchanged since that review, so I'm carrying all of them forward as still-open
below rather than treating anything as resolved.
**Issue 1: `clearCoordinatorService()` can close a newer master's listeners
if it races an in-flight `initCoordinatorService()` during `shutdown()`**
- Location:
`seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/CoordinatorService.java:1211-1213`
- Problem: see the runtime-flow trace above — the field read/write is
unsynchronized across the two methods, and `shutdown()` gives a narrow window
where a stale/late tick and a shutdown-triggered clear can interleave.
- Potential risk: the *new* active master's finished-job expiration side
effects silently stop firing cluster-wide until the next role switch, with no
error or log line explaining why — a worse failure mode than the leak this PR
fixes.
- Best improvement: apply the same capture-and-clear pattern already used
for `resourceManager` two lines below (a local variable captured under the same
`synchronized` block, or an `AtomicReference` CAS shared with
`initCoordinatorService()`), so `close()` can never target an instance
published by a later activation.
- Severity: Medium
- Raised by another reviewer: Yes (@SEZ9)
**Issue 2: closed instance left reachable via the field with no `closed`
marker**
- Location: `CoordinatorService.java:1211-1213`
- Problem: nothing distinguishes a closed `JobHistoryService` from a live
one at the type level.
- Potential risk: purely forward-looking — only bites if a future refactor
adds an instance-reuse branch to `initCoordinatorService()`; not reachable
today since that method unconditionally constructs a fresh instance every time.
- Best improvement: a `closed` boolean (cheap insurance), or fold into Issue
5's `AutoCloseable` change.
- Severity: Low
- Raised by another reviewer: Yes (@SEZ9, who rated it Medium — I rate Low
since it's not reachable in the current code path)
**Issue 3: constructor can still orphan 1-2 registrations if
`addEntryListener` throws mid-construction**
- Location:
`seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/master/JobHistoryService.java:152-160`
- Problem: if the second/third `addEntryListener` throws (e.g.
`HazelcastInstanceNotActiveException` racing a shutdown mid-activation), the
constructor aborts and the instance holding the earlier id(s) is never
published, so `close()` can never reach them.
- Potential risk: a narrow, rare trigger; strictly better than the pre-fix
baseline (guaranteed 3-listener leak on every transition) since at most 1-2
registrations are orphaned on an exceptional path.
- Best improvement: wrap the three registrations in try/catch, best-effort
removing any already-registered ids via `removeEntryListenerQuietly` before
rethrowing.
- Severity: Low
- Raised by another reviewer: Yes (@SEZ9, who rated it Medium — I rate Low
as it's a residual gap, not a regression, and narrowly triggered)
**Issue 4: `removeEntryListenerQuietly` treats the expected shutdown race
and a genuine bug identically, and discards the boolean result of
`removeEntryListener`**
- Location: `JobHistoryService.java:450-457`
- Problem: catches bare `Exception` at a single `warning` level; also
discards the `false` return that `removeEntryListener` gives when a
registration is already gone, so a persistent, non-exception removal failure
would silently reintroduce the exact leak this PR targets.
- Best improvement: catch `HazelcastInstanceNotActiveException` separately
at fine/debug level, log other exceptions at warning with the `registrationId`,
consider logging on an unexpected `false` return.
- Severity: Low
- Raised by another reviewer: Yes (@SEZ9)
**Issue 5: `JobHistoryService` doesn't declare `implements AutoCloseable`**
- Location: `JobHistoryService.java:69`
- Problem: the class now owns a `close()` lifecycle matching
`AutoCloseable`'s signature exactly, but doesn't declare it.
- Best improvement: declare `implements AutoCloseable` — free, makes the
"must be closed" contract compiler/IDE-discoverable, and lets the new test use
try-with-resources (also fixing Issue 7 below for free).
- Severity: Low
- Raised by another reviewer: Yes (@SEZ9)
**Issue 6: test-only accessor `getEntryListenerRegistrationIds()` isn't
marked `@VisibleForTesting`**
- Location: `JobHistoryService.java:202`
- Problem: package-private, Javadoc says "only intended for tests," but
nothing enforces that; `@VisibleForTesting` is already an established
convention in this exact module.
- Best improvement: add the annotation.
- Severity: Low
- Raised by another reviewer: Yes (@SEZ9)
**Issue 7: `testCloseRemovesFinishedJobEntryListeners`'s `liveService` (and
the create/close loop) register listeners on the shared test node without
try/finally**
- Location:
`seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/master/JobHistoryServiceListenerCleanupTest.java:77-104`
- Problem: I confirmed directly — line 77's `liveService` has no surrounding
try/finally (the file's only try/finally block is elsewhere, guarding the
isolated-node test). A failed assertion between creation and cleanup leaves
live registrations on the shared node's cluster-wide IMaps for the rest of the
JVM's test run, reproducing this PR's exact leak pattern inside the test suite
itself and making unrelated failures in sibling tests harder to diagnose.
- Best improvement: wrap `liveService` and each loop-created instance in
try/finally calling `close()` (documented idempotent, so double cleanup is
safe).
- Severity: Low
- Raised by another reviewer: Yes (@SEZ9)
## 1.2 Compatibility Impact
Fully compatible. Purely additive: three new `private final UUID` fields,
one new `public void close()`, one new package-private test accessor, and one
new call site in an existing method. No config option, public API signature,
checkpoint/savepoint format, or wire-protocol change. `close()` is idempotent
and cannot throw into the caller.
## 1.3 Performance / Side-Effect Analysis
Negligible. `close()` does at most three `IMap.removeEntryListener` calls
once per master-role transition — transitions are rare relative to normal
operation, not per-job or per-event. No new locks, threads, or hot-path impact.
## 1.4 Error Handling and Logging
Adequate for the stated best-effort contract, with the diagnosability gap
noted as Issue 4 above (bare `Exception` catch, discarded boolean result).
# 2. Code Quality Assessment
## 2.1 Coding Standards
New fields/methods and the test class carry Javadoc explaining purpose and
lifecycle risk. No missing-comment gaps on the core methods. Issues 5/6 above
are the only structural/annotation gaps.
## 2.2 Test Coverage and Test Stability
Two new regression tests cover both the direct `close()` path and the
production `clearCoordinatorService()` call path, using
`Awaitility.await().ignoreExceptions().atMost(30,
TimeUnit.SECONDS).untilAsserted(...)` for the isolated-node activation wait —
no `Thread.sleep`, no fixed-delay guessing. The isolated second Hazelcast node
uses a reserved port range (`TestUtils.getAvailablePort(100)`) rather than a
default/shared port, avoiding cross-test-class port collisions in the same JVM.
**Rating: Stable**, with the one test-hygiene caveat in Issue 7 (no try/finally
on the shared-node registrations).
## 2.3 Documentation Updates
Not applicable — no user-facing config, API, or behavior contract changes;
the PR description correctly marks no documentation update as required.
# 3. Architectural Soundness
## 3.1 Elegance of the Solution
Precise fix for the reported leak on its primary (scheduler-driven) path.
Minimal and additive; follows the file's own established idioms almost
everywhere except the one place that matters most for correctness under the
shutdown race (Issue 1).
## 3.2 Maintainability
Good — `close()` and the new fields make the listener-ownership contract
explicit where none existed before. Issues 2/5/6 are all about making that
contract harder to violate accidentally in the future.
## 3.3 Extensibility
No concerns; the pattern generalizes cleanly if a fourth cluster-wide
listener is ever added.
## 3.4 Historical-Version Compatibility
No serialized state, checkpoint format, or config option touched. A rolling
upgrade with old and new binaries mixed is safe: old nodes keep leaking as
before until upgraded, new nodes stop leaking, and there is no cross-version
wire-format interaction.
# 4. Issue Summary
| # | Issue | Location | Severity |
|---|---|---|---|
| 1 | `clearCoordinatorService()` can close a newer master's listeners if it
races an in-flight `initCoordinatorService()` during `shutdown()` |
`CoordinatorService.java:1211-1213` | Medium |
| 2 | Field left pointing at a closed instance with no `closed` marker
(forward-looking, not an active bug today) |
`CoordinatorService.java:1211-1213` | Low |
| 3 | Partial-construction failure can orphan 1-2 listener registrations
with no owner (not a regression) | `JobHistoryService.java:152-160` | Low |
| 4 | `removeEntryListenerQuietly` treats expected-shutdown and genuine
failures identically, discards `removeEntryListener`'s boolean result |
`JobHistoryService.java:450-457` | Low |
| 5 | `JobHistoryService` should declare `implements AutoCloseable` |
`JobHistoryService.java:69` | Low |
| 6 | Test-only accessor should be `@VisibleForTesting` |
`JobHistoryService.java:202` | Low |
| 7 | `testCloseRemovesFinishedJobEntryListeners` lacks try/finally cleanup
on the shared test node | `JobHistoryServiceListenerCleanupTest.java:77-104` |
Low |
CI note: `Build` was red on this head as of yesterday's check; I traced the
two failing jobs back then (`unit-test (8, windows-latest)` — Maven
dependency-resolution connection reset for `seatunnel-transforms-v2`, before
any engine-server test runs; `all-connectors-it-1 (ubuntu-latest)` — the known
`CouchbaseIT` container-startup flake that's been hitting `dev` since
2026-08-15) and both are infra-level, unrelated to this diff. Nothing has
changed on the CI front since then that I can see from this pass; I'd still
recommend rerunning the two failed jobs for a clean signal rather than treating
the red `Build` as evidence against this PR.
# 5. Merge Recommendation
### Conclusion: Ready to merge after fixes
1. Blockers — must be fixed
- **Issue 1** (Medium): apply the capture-and-clear idiom already used
for `resourceManager` right below, to `jobHistoryService`, in
`clearCoordinatorService()`, so `close()` can never target an instance
published by a later `initCoordinatorService()` call. This is a small,
mechanical change and it protects the exact correctness guarantee this PR
exists to establish.
2. Recommended fixes — non-blocking
- Issues 2-7 (closed-instance marker / `AutoCloseable`,
`@VisibleForTesting`, partial-construction cleanup, richer removal-failure
logging, test try/finally hygiene).
Overall assessment: unchanged from yesterday's round — the core fix is
correct and well-tested for the normal, scheduler-driven master-transition path
issue #11807 describes. `@SEZ9`'s findings remain legitimate and unaddressed
since no new commit has landed; Issue 1 in particular is cheap enough to fix
now that I'd rather land it in this same PR than merge with a known (if narrow)
concurrency hazard in code whose entire purpose is closing a
lifecycle-correctness gap. Issues 2-7 are genuine but low-cost hardening that
can reasonably follow in this PR or a fast follow-on without blocking the
primary fix much longer.
--
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]