DanielLeens opened a new pull request, #11809:
URL: https://github.com/apache/seatunnel/pull/11809

   ### Purpose of this pull request
   
   Fixes #11807.
   
   `JobHistoryService` registers three expiration listeners on the cluster-wide 
finished-job IMaps (`IMAP_FINISHED_JOB_STATE`, `IMAP_FINISHED_JOB_METRICS`, 
`IMAP_FINISHED_JOB_VERTEX_INFO`) in its constructor, and 
`CoordinatorService.checkNewActiveMaster()` creates a new `JobHistoryService` 
instance every time a node becomes the active master. Before this PR there was 
no deregistration path: the registration ids returned by 
`addEntryListener(...)` were dropped, and `clearCoordinatorService()` closed 
the resource manager and the event processor but never touched the listeners.
   
   As a result, every active -> inactive -> active master transition (master 
failover, split-brain healing, rolling restart scenarios) leaked the previous 
three listeners:
   
   * the stale registrations keep the old `JobHistoryService` instance 
reachable from the IMap listener registry (the listeners are non-static inner 
classes), retaining the old service objects in heap;
   * expiration side effects run once per leaked instance, so a single expired 
`JobDAGInfo` entry triggers duplicate `CleanLogOperation` fan-out to member 
nodes and over-counts the finished-job cleanup totals.
   
   Changes, following the fix direction proposed in the issue:
   
   1. `JobHistoryService` now stores the three listener registration ids 
returned by `addEntryListener(...)`.
   2. New `JobHistoryService.close()` deregisters exactly those listeners. 
Removal is best effort per map, so a Hazelcast instance that is already 
shutting down cannot break the master switch flow, and calling close twice is a 
safe no-op.
   3. `CoordinatorService.clearCoordinatorService()` calls 
`jobHistoryService.close()` when the node leaves the active master role (this 
also covers coordinator shutdown, which goes through the same method). The 
`jobHistoryService` field itself is intentionally kept non-null, so read paths 
that still hold the old instance (for example REST handlers on a node that just 
lost the master role) behave exactly as before; only the listener registrations 
are removed.
   
   The normal single-master steady state is not affected: `close()` only runs 
when a node leaves the active master role or shuts down, and the new active 
master registers fresh listeners through the new instance as before.
   
   ### Does this PR introduce _any_ user-facing change?
   
   No. No config options, public API contracts, or serialized formats change. 
The only behavior difference is that after a master role switch, finished-job 
expiration side effects no longer fire additionally from stale listeners on the 
node that left the master role, which is the bug being fixed.
   
   ### How was this patch tested?
   
   Added `JobHistoryServiceListenerCleanupTest` with two regression tests:
   
   * `testCloseRemovesFinishedJobEntryListeners`: a positive control first 
proves the exposed registration ids are live registrations 
(`removeEntryListener` returns true for a service that is not closed), then 
three create/close cycles emulate repeated active-master transitions and assert 
that after `close()` every registration id is gone (`removeEntryListener` 
returns false).
   * `testClearCoordinatorServiceDeregistersJobHistoryListeners`: starts an 
isolated SeaTunnel Hazelcast instance, captures the registration ids of the 
active coordinator's `JobHistoryService`, calls `clearCoordinatorService()`, 
and asserts the captured registrations are deregistered through the production 
cleanup path.
   
   Both tests only use the public Hazelcast `IMap` listener API for their 
assertions. Compilation and the full test suite run through GitHub CI for this 
PR; locally only Spotless formatting was applied.
   
   ### Check list
   
   * [x] No new jar or binary dependency is added.
   * [x] No documentation update is required: no user-facing option or behavior 
contract changes.
   * [x] No entry in `incompatible-changes.md` is required: the change is 
backward compatible.
   * [x] Not a connector contribution, so the connector file checklist does not 
apply.
   


-- 
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]

Reply via email to