DanielLeens commented on PR #11856: URL: https://github.com/apache/seatunnel/pull/11856#issuecomment-5739315484
Hi @SEZ9, thank you for coming back to this thread, and for confirming that `dc2c64103c` closes the compile break. This is my own PR, so I am answering as the author and have re-checked everything below on the current head `dc2c64103c` (unchanged since my last comment) and in the fork CI run for that head. ## Direct answers **F1 - where the listener is registered, where the write happens, and whether a real lifecycle path is tested** - Registration: `SeaTunnelServer.init()` calls `getLifecycleService().addLifecycleListener(this)` at `SeaTunnelServer.java:154` (the class implements `LifecycleListener`). - Dispatch: `stateChanged(...)` at `SeaTunnelServer.java:239-245` sends `SHUTTING_DOWN` to `markLocalGracefulMemberRemoval()` (`:284-286`) and `STARTED` to the clear (`:292-294`). - The IMap write is `gracefulMemberRemovalIMap.put(thisAddress, System.currentTimeMillis(), TTL, MILLISECONDS)` at `SeaTunnelServer.java:309-313`, inside `updateLocalGracefulMemberRemovalMarker` (`:300-326`). `ManagedService.shutdown()` (`:247-278`) no longer touches the map, and `SeaTunnelServerShutdownTest.shouldNotMarkGracefulMemberRemovalDuringManagedServiceShutdown` pins that. - Real lifecycle test: `ClusterFaultToleranceIT.testGracefulShutdownPublishesMemberRemovalMarker` (`ClusterFaultToleranceIT.java:79-146`) starts two real members with `SeaTunnelServerStarter.createHazelcastInstance`, registers an entry listener on node2's real map proxy (`:108-119`), calls `node1.shutdown()` (`:121`), and asserts the marker arrives and is accepted by `CoordinatorService.isGracefulMemberRemovalMarkerValid` (`:124-131`). The mock-based `SeaTunnelServerShutdownTest` only covers the branch logic. - Evidence that it ran: in fork run `35240527669` (this head, attempt 3), the `engine-v2-it (8, ubuntu-latest)` log reports `Tests run: 10, Failures: 0, Errors: 0, Skipped: 2 - in ClusterFaultToleranceIT`. The two skips match the two `@Disabled` tests that already exist in that file on `dev` and on this head. That the new test itself passed is an inference from the class-level result; I did not find a per-method line for it. - What that test does not cover: it exercises `HazelcastInstanceImpl.shutdown()`, not the JVM shutdown-hook thread (SIGTERM) path. That the hook reaches the same `LifecycleServiceImpl.shutdown` and fires `SHUTTING_DOWN` is from my reading of the Hazelcast 5.1 sources, not from a test in this PR. **F4 - TTL versus timestamp** Both, with different jobs. Eviction is a Hazelcast-side per-entry TTL: the 4-argument `put(key, value, ttl, unit)` at `SeaTunnelServer.java:309-313` with `Constant.GRACEFUL_MEMBER_REMOVAL_MARK_TTL_MILLIS` = 5 minutes (`Constant.java:72`). Separately, `CoordinatorService.isGracefulMemberRemovalMarkerValid` (`CoordinatorService.java:2157-2161`) checks `abs(now - markedAt) <= 5 min`. That second check does compare the departing member's `System.currentTimeMillis()` with the coordinator's clock, so it tolerates up to 5 minutes of skew; beyond that the marker is rejected and the failure stays at `ERROR`, which is the safe direction. So the "no eviction" part of F4 is resolved by the native TTL, and the timestamp is only a secondary guard. **F2 - destructive removal and failover** The read is non-destructive: `gracefulMemberRemovalIMap.get(lostAddress)` at `CoordinatorService.java:2185`. The clear runs only after the propagation loop (`:2222-2238`), and only when no master-switch restore is in flight (`canClearGracefulMemberRemovalMarker`, `:2169-2172`, called at `:2239-2245`). It is a value-conditional `remove(lostAddress, markedAt)` (`:2207`), so it cannot delete a newer marker from a restarted member at the same address. When the clear is skipped, or a coordinator dies mid-processing, the entry stays in the IMap until its TTL, and the new master classifies through `PhysicalVertex.checkTaskGroupIsExecuting` (`PhysicalVertex.java:239-250`), which reads the same marker without clearing it (`:278-288`) and applies the same helpers (`recordMemberRemovedFailure`, `:760-770`). Unit tests: `shouldRetainGracefulMemberRemovalMarkerDuringMasterSwitchRecovery` and `shouldClassifyMissingWorkerDuringMasterFailoverAsGraceful`. One nit I noticed while re-reading: the Ja vadoc at `PhysicalVertex.java:273-277` still says "one-time" marker, which is stale wording now that the read is non-destructive. ## Per-finding status on `dc2c64103c` | Finding | Status | Evidence | |---|---|---| | F1 marker written while PASSIVE | Fixed in source; publication exercised by a real two-node IT that passed on this head | `SeaTunnelServer.java:154`, `:239-245`, `:309-313`; `ClusterFaultToleranceIT.java:79-146` | | F2 consumed before processing / failover | Fixed | `CoordinatorService.java:2185`, `:2207`, `:2239-2245`; `PhysicalVertex.java:278-288` | | F3/F6 `JobException` replaced by a string | Fixed; payload shape unchanged | `buildMemberRemovedFailureState`, `CoordinatorService.java:2144-2149`; `shouldKeepThrowablePayloadForMemberRemovedFailureState` | | F4 eviction | Fixed by native TTL; timestamp is a secondary guard | `SeaTunnelServer.java:309-313`; `CoordinatorService.java:2157-2161` | | F5/F8 message matching and duplicated template | Fixed | no `Pattern` or `.matches(` in `PhysicalVertex.java`; single producer `buildMemberRemovedOfflineMessage` (`CoordinatorService.java:2133-2137`), called from `:2146` and `PhysicalVertex.java:768`; the classification is a typed boolean | | F7 docs | Present | `docs/en/engines/zeta/state-storage-and-recovery.md` (new section at line 155) and the zh counterpart | ## What CI actually shows on this head I read fork run `35240527669` (attempt 3, finished 2026-09-18T14:13Z; the apache-side `Build` check finished as failure at 14:14Z). It is red, so I cannot call this PR verified or ready to merge yet. - The compile break is fixed: `unit-test` passed on Java 8 and 11 on both ubuntu and windows, and `engine-v2-it (11)` passed. 79 of 94 jobs succeeded, 11 were skipped, 3 failed, and 1 was cancelled. - `engine-v2-it (8)` failed with 204 tests run, 1 failure and 1 error: `BackpressureSlowSinkIT.testCheckpointsKeepCompletingUnderSustainedBackpressure` (only 1 of the expected 3 additional checkpoints), and `SplitClusterFaultToleranceIT.testStreamJobCancelResolvesWhenWorkerCrashesBeforeCancelAck` (`expected: <CANCELED> but was: <FAILED>`). Both match the targets of open PRs #12316 and #12311, so a `dev` sync would not clear them. My diff does not touch the vertex-state filter or the FAILED target in `makeTasksFailed` (`CoordinatorService.java:2256-2265`); it only adds the classification argument. That these two are unrelated to my change is my inference, not something a run has proven. - `all-connectors-it-2 (8)` and `(11)` both failed on `OpengaussCDCIT.testAddFieldWithRestore` (`ConditionTimeout` at `:476`). This is a CDC connector test unrelated to the engine shutdown path, as far as I can tell. - `jdbc-connectors-it-part-1 (11)` was cancelled at 14:13Z after running about two hours, with no failed test in its log. I have not established why it was cancelled. ## Runtime evidence for the WARN downgrade (the acceptance list from my last comment) Counts below are from the `engine-v2-it (8)` log only, which is one job of the run. - `... end with state FAILED due to node offline` at `WARN`: 125 occurrences (baseline on the previous head: 0). No `ERROR`-level offline line remains in that log. - `NoDataMemberInClusterException`: 0 occurrences. - `Failed to clear graceful member removal marker`: 13, all `HazelcastInstanceNotActiveException`, logged from the coordinator's clear (`CoordinatorService.java:2207`). This is not fully gone; it fires when the coordinator itself is stopping. - `Failed to mark graceful member removal marker`: 39, all `HazelcastInstanceNotActiveException`, so the criterion "gone for shutdown() paths" is not literally met. I traced four of them in `ClusterFaultToleranceIT`, plus one in `SplitClusterPendingJobLifecycleFailoverIT`, and they are calls to `shutdown()` on an instance that was already down: for example `testBatchJobRestoreIn2NodeWorkerDown` shuts node2 down in the body (`:441`) and again in `finally` (`:470`), and the same pattern holds at `:540`/`:591`, `:660`/`:690` and `:763`/`:809`. `LifecycleServiceImpl.shutdown` fires `SHUTTING_DOWN` on every call, so the listener tries a write against a dead instance. I have not traced the other 34 individually, so I cannot claim they are all of this kind. ## Remaining gaps, stated plainly 1. CI on this head is red (list above). I am not treating the PR as ready until the failed jobs are re-run and pass, or are shown to be unrelated. 2. No test asserts the `WARN` versus `ERROR` classification in a real cluster. The engine-level evidence is the log counts above; the classification rules themselves are covered by unit tests only (`PhysicalVertexTest`, `CoordinatorServiceMemberRemovedTest`). 3. My earlier recommendation N2 is still open on this head: the marker `put` (`SeaTunnelServer.java:309`) and the `STARTED` clear (`:315`) are synchronous calls on lifecycle threads. Together with the stack-trace noise above, a bounded async call, or skipping the write when the instance is already not running, would be a reasonable follow-up. 4. N4 (one sentence in the docs about the 5-minute TTL versus `hazelcast.max.no.heartbeat.seconds`) is also still open. I will do a final pass on the next green run of this head. Thanks again for the careful review. -- 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]
