DanielLeens commented on PR #11856:
URL: https://github.com/apache/seatunnel/pull/11856#issuecomment-5701028678

   [REVIEW] Maintainer re-review of the current head `fd60dfd596` (a pure `dev` 
merge on top of `e6e243744e`; I verified the three-dot diff against `dev` is 
content-identical to the previously reviewed head, 31 files, +1148/-42).
   
   **Verdict: not ready to merge. Two blocking findings, both confirmed from 
the real Hazelcast 5.1 source and from the engine E2E logs of the last green CI 
run.** This retracts the "Ready to merge" conclusion I posted on 2026-09-06 at 
`e6e243744e`. That conclusion rested on green CI plus source-level tracing of 
the `GracefulShutdownAwareService` hook. It never checked what the engine E2E 
logs actually recorded on the graceful-shutdown path. I did that this round, 
and the logs contradict the earlier conclusion: the WARN downgrade has never 
fired in CI, and the marker write fails on every graceful shutdown.
   
   # What this PR is supposed to do
   
   When a worker leaves the cluster during an intentional scale-down, 
`CoordinatorService.failedTaskOnMemberRemoved` fails its task groups and 
`PhysicalVertex.stateProcess()` logs that at `ERROR`. The PR has a node that is 
shutting down gracefully write a TTL-bounded marker into the new 
`engine_gracefulMemberRemoval` IMap, and the coordinator downgrades only a 
marker-proven departure to `WARN`. The PR body says the change "keeps the 
existing failure state handling unchanged and only narrows the log level for 
the known node-offline case".
   
   # Runtime chain re-traced on this head
   
   SeaTunnel side is the PR head; Hazelcast side is the vanilla 5.1 sources 
jar. The shaded `seatunnel-shade-hazelcast:5.1-3.0.0` jar that CI runs carries 
byte-identical `Node.class` and `Invocation.class` (md5 checked), so the 
vanilla source is authoritative here.
   
   ```text
   Graceful shutdown (intended marker write)
     HazelcastInstance.shutdown()  /  NodeShutdownHookThread with policy 
GRACEFUL   Node.java:751-779
       -> LifecycleServiceImpl.shutdown(terminate=false)                        
    LifecycleServiceImpl.java:92-105
            fireLifecycleEvent(SHUTTING_DOWN)   (synchronous, node still 
ACTIVE)    :94, :66-72
            Node.shutdown(false)                                                
    Node.java:502
              setShuttingDown(): shuttingDown=true, state=PASSIVE               
    Node.java:633-635
              callGracefulShutdownAwareServices()  (graceful-shutdown executor) 
    Node.java:520, 554-579
                -> SeaTunnelServer.onShutdown()                                 
    SeaTunnelServer.java:242
                  -> updateLocalGracefulMemberRemovalMarker(true)               
    SeaTunnelServer.java:310-343
                    -> IMap.put(thisAddress, now, TTL)                          
    SeaTunnelServer.java:319-323
                      -> MapProxyImpl.put -> AbstractDistributedObject.toData()
                        -> getNodeEngine() -> lifecycleCheck()                  
    AbstractDistributedObject.java:104-115
                          -> NodeEngineImpl.isRunning() -> Node.isRunning()     
    NodeEngineImpl.java:394, Node.java:646-648
                             == !shuttingDown.get() == false
                          => HazelcastInstanceNotActiveException                
    caught at SeaTunnelServer.java:335
              shutdownServices() -> SeaTunnelServer.shutdown(false)
   
   Coordinator classification
     MembershipManager.sendMembershipEventNotifications -> membership-event 
executor   MembershipManager.java:819-833
       -> SeaTunnelServer.memberRemoved  (catches SeaTunnelEngineException 
only)      SeaTunnelServer.java:349-357
         -> CoordinatorService.memberRemoved                                    
      CoordinatorService.java:2249-2254
           -> failedTaskOnMemberRemoved                                         
      CoordinatorService.java:2197-2226
             -> getGracefulMemberRemovalMarker -> IMap.get   (no try/catch)     
      CoordinatorService.java:2179-2184
             -> makeTasksFailed(...)                                            
      CoordinatorService.java:2228-2246
             -> clearGracefulMemberRemovalMarker -> IMap.remove (no try/catch)  
      CoordinatorService.java:2191-2195
   
   Log level
     PhysicalVertex.stateProcess case FAILED -> shouldLogFailureAsWarn          
      PhysicalVertex.java:666-692, 728-730
   ```
   
   # Findings
   
   ## Issue 1 (P0, blocking): the marker write is unreachable on the real 
graceful-shutdown path. The WARN downgrade never fires, and every graceful 
shutdown now logs a new WARN with a stack trace.
   
   **Trigger condition:** any graceful shutdown of a SeaTunnel node 
(`HazelcastInstance.shutdown()`, or SIGTERM under the new `GRACEFUL` hook 
policy). This is the exact scenario the PR targets.
   
   **Root cause (Hazelcast 5.1 source):**
   - `Node.shutdown()` flips `shuttingDown` to `true` in `setShuttingDown()` 
(`Node.java:633-635`) before it calls `callGracefulShutdownAwareServices()` 
(`Node.java:520`).
   - Every `IMap` proxy method goes through 
`AbstractDistributedObject.getNodeEngine()`, whose `lifecycleCheck()` throws 
`HazelcastInstanceNotActiveException` when `!engine.isRunning()` 
(`AbstractDistributedObject.java:104-115`). `NodeEngineImpl.isRunning()` 
delegates to `Node.isRunning()`, which is `!shuttingDown.get()` 
(`NodeEngineImpl.java:394`, `Node.java:646-648`).
   - Therefore no IMap proxy operation can succeed from inside 
`GracefulShutdownAwareService.onShutdown()`, deterministically, regardless of 
node state, cluster state, or partition ownership. The earlier rounds' claim 
that this hook runs "while map and operation services are still active" was 
wrong: the operation service is still running, but the proxy layer rejects the 
call.
   
   **Empirical evidence (fork run `33966670706`, head `e6e243744e`, whose diff 
is content-identical to this head):**
   
   | counter in the engine E2E job log | `engine-v2-it (8)` | `engine-v2-it 
(11)` |
   |---|---|---|
   | `Failed to mark graceful member removal marker` (SeaTunnelServer WARN) | 
153 | 152 |
   | of which followed by `HazelcastInstanceNotActiveException: Hazelcast 
instance is not active!` | 139 (rest are interleaved log lines) | 139 |
   | `Graceful shutdown failed for 
org.apache.seatunnel.engine.server.SeaTunnelServer` (Hazelcast WARN) | 153 | 
152 |
   | `end with state FAILED due to node offline` (the new WARN branch) | **0** 
| **0** |
   | `end with state FAILED and Exception: ... deployed node(...) offline` (the 
pre-PR ERROR branch) | 122 | 124 |
   
   Sample stack from the JDK 8 log (`ClusterFaultToleranceIT`-family tests call 
`node.shutdown()` on real multi-node clusters):
   ```text
   WARN  org.apache.seatunnel.engine.server.SeaTunnelServer - Failed to mark 
graceful member removal marker
   com.hazelcast.core.HazelcastInstanceNotActiveException: Hazelcast instance 
is not active!
        at 
com.hazelcast.spi.impl.AbstractDistributedObject.throwNotActiveException(AbstractDistributedObject.java:115)
        at 
com.hazelcast.spi.impl.AbstractDistributedObject.lifecycleCheck(AbstractDistributedObject.java:110)
        at 
com.hazelcast.spi.impl.AbstractDistributedObject.getNodeEngine(AbstractDistributedObject.java:104)
        at 
com.hazelcast.spi.impl.AbstractDistributedObject.toData(AbstractDistributedObject.java:78)
        at com.hazelcast.map.impl.proxy.MapProxyImpl.put(MapProxyImpl.java:137)
        at 
org.apache.seatunnel.engine.server.SeaTunnelServer.updateLocalGracefulMemberRemovalMarker(SeaTunnelServer.java:319)
        at 
org.apache.seatunnel.engine.server.SeaTunnelServer.markLocalGracefulMemberRemoval(SeaTunnelServer.java:296)
        at 
org.apache.seatunnel.engine.server.SeaTunnelServer.onShutdown(SeaTunnelServer.java:243)
        at com.hazelcast.instance.impl.Node$2.run(Node.java:564)
   WARN  com.hazelcast.instance.impl.Node - [localhost]:5801 [...] Graceful 
shutdown failed for org.apache.seatunnel.engine.server.SeaTunnelServer@362203b
   ```
   
   **Impact:** the feature is inert in production. Classification stays at 
`ERROR` exactly as on `dev`, and on top of that every graceful shutdown emits 
two new WARN lines, one with a full stack trace. For a PR whose purpose is to 
reduce misleading log noise, this is a net regression. Green CI proves nothing 
here: the failure is caught and logged at WARN, so no test fails.
   
   **Why the tests did not catch it:** 
`SeaTunnelServerShutdownTest.shouldMarkGracefulMemberRemovalOnGracefulShutdown` 
mocks `IMap`, so by construction it cannot observe the proxy's lifecycle check. 
This is the second time on this PR that the marker was published from a hook 
where the write cannot succeed (SEZ9's original finding was 
`ManagedService.shutdown()`); a mock-only test cannot guard SPI placement.
   
   **Fix direction (minimal):** publish the marker from 
`LifecycleListener.stateChanged(SHUTTING_DOWN)`. 
`LifecycleServiceImpl.shutdown()` fires that event synchronously at 
`LifecycleServiceImpl.java:94`, before `node.shutdown()` at `:101`, so 
`node.isRunning()` is still `true` and node state is `ACTIVE`; a plain 
`IMap.put` works there. `SeaTunnelServer` is already registered as a 
`LifecycleListener` at `SeaTunnelServer.java:156`, so this removes the 
`GracefulShutdownAwareService` implementation rather than adding code. 
Trade-off to state explicitly in the PR: `SHUTTING_DOWN` fires for both 
`shutdown()` and `terminate()`. Both are intentional in-process exits (an API 
call or the JVM shutdown hook); crashes, `kill -9`, OOM kills and network 
partitions never fire it, so the "intentional versus unproven" distinction the 
PR cares about is preserved. If forced `terminate()` must still be excluded, 
that needs a different signal, and whatever is chosen must be proven with a 
real-cluster tes
 t (Issue 7).
   
   ## Issue 2 (P0, blocking): a new unguarded distributed `IMap.get` at the 
head of `failedTaskOnMemberRemoved` can abort task-failure propagation for an 
entire lost member.
   
   **Trigger condition:** any member removal (this is the normal path for every 
worker loss, graceful or not) where 
`gracefulMemberRemovalIMap.get(lostAddress)` throws.
   
   **Root cause:** `failedTaskOnMemberRemoved` 
(`CoordinatorService.java:2197-2226`) now calls 
`getGracefulMemberRemovalMarker` (`:2179-2184`), a synchronous distributed 
`IMap.get`, before the `runningJobMasterMap.forEach(...)` loop, with no 
exception handling. The exception propagates through 
`CoordinatorService.memberRemoved` (`:2249-2254`) to 
`SeaTunnelServer.memberRemoved` (`SeaTunnelServer.java:349-357`), which only 
catches `SeaTunnelEngineException`, and then to the Hazelcast membership-event 
executor (`MembershipManager.java:831`), where it is swallowed as an 
executor-level log. Result: no task group on the lost worker is marked 
`FAILED`, the pipeline never restarts, and the job sits in `RUNNING` with tasks 
that no longer exist. There is no other detector for this: 
`checkTaskGroupIsExecuting` only runs on master-failover restore and on cancel.
   
   **Why this is new:** on `dev`, everything before the per-vertex 
`updateTaskState` is in-memory (`getCurrentExecutionAddress` reads the job 
master's owned-slot map, `getExecutionState()` returns the `currExecutionState` 
field), and `updateTaskState` wraps its IMap access in 
`RetryUtils.retryWithException` plus a per-vertex `catch (Exception e)` 
(`PhysicalVertex.java:390-446`). So on `dev` a Hazelcast hiccup costs at most 
one vertex; on this head it costs every vertex of every job on the lost member.
   
   **When `get` throws in practice:** right at member removal, Hazelcast is 
repartitioning the partitions the lost member owned. A `get` that lands on such 
a partition retries on `WrongTargetException`/`TargetNotMemberException` up to 
`hazelcast.invocation.max.retry.count`, which the shipped configs set to `20`, 
with the backoff in `Invocation.handleRetry` (`Invocation.java:706-712`): 
roughly 3.5 s of budget before the exception propagates. Slow owners hit 
`OperationTimeoutException` at the 60 s call timeout instead. At the load 
profile this engine is reviewed against (trillions of rows per day, tens of 
billions per job, 271 partitions repartitioning under checkpoint state 
pressure), a >3.5 s repartition at the moment a worker disappears is a 
plausible production event, not a corner case.
   
   **Fix direction:** wrap the marker read (and the trailing 
`clearGracefulMemberRemovalMarker` at `:2222-2225`) in `try/catch` and treat 
any failure as "unproven" (`null` marker, `ERROR` path). 
`PhysicalVertex.getGracefulMemberRemovalMarker` (`PhysicalVertex.java:278-288`) 
already does exactly this on the restore path; the coordinator path must be at 
least as defensive, because it is the primary path.
   
   ## Issue 3 (Medium): `hazelcast.shutdownhook.policy: GRACEFUL` changes 
SIGTERM semantics for every shipped deployment, is unrelated to logging, and is 
unnecessary once Issue 1 is fixed correctly.
   
   - Default policy is `TERMINATE` (`ClusterProperty.java:1533-1534`). The PR 
sets `GRACEFUL` in 3 shipped `config/*.yaml`, the engine-common default 
`hazelcast.yaml`, 2 Kubernetes conf files and 11 docs pages.
   - `stop-seatunnel-cluster.sh` sends a plain `kill` (SIGTERM) at `:51` and 
`:58`. Under `GRACEFUL`, the JVM hook now runs `lifecycleService.shutdown()`, 
so `InternalPartitionServiceImpl.onShutdown` 
(`InternalPartitionServiceImpl.java:897-935`) blocks until the master has 
migrated this member's partitions away, bounded by 
`hazelcast.graceful.shutdown.max.wait` = 600 s (`ClusterProperty.java:627-628`, 
not overridden anywhere in the repo). Stopping a data member (master, or a 
hybrid MASTER_AND_WORKER node) with a large IMap footprint now keeps the JVM 
alive noticeably longer; on Kubernetes the pod can be SIGKILLed at 
`terminationGracePeriodSeconds` mid-migration.
   - None of this is described in the PR body or in the added docs, which only 
say operators "must retain the GRACEFUL policy". A change to what SIGTERM does 
on every node needs its own justification, its own doc paragraph about the wait 
and its timeout, and a release-note entry.
   - With the Issue 1 fix on `SHUTTING_DOWN`, the marker is written under the 
default `TERMINATE` policy as well (`terminate()` also goes through 
`LifecycleServiceImpl.shutdown(true)` and fires `SHUTTING_DOWN` at `:94`). I 
recommend dropping the 18 config/doc/Kubernetes edits from this PR, or 
splitting them into a dedicated PR with the rationale above.
   
   ## Issue 4 (Medium): the startup-side marker clear fails deterministically 
on lite-member workers and logs a WARN with a stack trace on every worker boot.
   
   - `stateChanged(STARTING)` (`SeaTunnelServer.java:252-256`) issues 
`removeAsync` from inside `Node.start()` at `Node.java:460`, before `join()` at 
`:483`. Separated-cluster workers are lite members 
(`SeaTunnelServerStarter.java:101`), so pre-join the only known member is a 
lite member and `PartitionInvocation.newTargetNullException` throws 
`NoDataMemberInClusterException`, which extends `HazelcastException` and is not 
retryable (`NoDataMemberInClusterException.java:24`). The `.exceptionally` 
handler then logs a WARN with the full stack.
   - Evidence: 35 `Failed to clear graceful member removal marker` in each 
engine E2E job (22 `NoDataMemberInClusterException`, 8 `WrongTargetException` 
after retries were exhausted on master nodes, 4 
`HazelcastInstanceNotActiveException`) across 200 member starts.
   - Functional impact is bounded (Hazelcast TTL plus the coordinator's 
post-classification clear), but this is another new per-boot WARN with a stack 
trace from a noise-reduction PR.
   - Fix direction: fire the clear from a post-join lifecycle point. 
`LifecycleState.STARTED` is fired from 
`HazelcastInstanceFactory.constructHazelcastInstance` 
(`HazelcastInstanceFactory.java:241`) only after `Node.start()` has returned 
and the join completed, so the partition table has a data member to target. 
Keep it asynchronous and fail-safe as it is now.
   
   ## Issue 5 (Low): the cancel path now records an "offline" failure message 
where none existed before.
   
   `checkTaskGroupIsExecuting` (`PhysicalVertex.java:228-271`) now installs 
`The taskGroup(...) deployed node(...) offline` into the failure slot whenever 
the worker is missing. It is called not only from the restore path but also 
from `noticeTaskExecutionServiceCancel` (`PhysicalVertex.java:457`), so a 
`CANCELED` task whose worker has already left now completes its 
`TaskExecutionState` with that message (`:658-664`) instead of `null`. Probably 
harmless, but it contradicts "failure state handling unchanged" and should be 
listed in the PR description as a deliberate change.
   
   ## Issue 6 (Low): Javadoc and docs assert the opposite of the observed 
ordering.
   
   `SeaTunnelServer.onShutdown` Javadoc (`SeaTunnelServer.java:236-240`) and 
`docs/en|zh/engines/zeta/state-storage-and-recovery.md` state that the callback 
runs "while map and operation services are still active" and "writes each 
server member's address into this map before member removal". Issue 1 shows the 
IMap proxy rejects the write there. Both must be rewritten with the fix so the 
next reader is not misled the way this review thread was.
   
   ## Issue 7 (Low on its own, required with the fix): no real-cluster test 
covers the PR's central claim.
   
   All three new test classes are Mockito-based. `ClusterFaultToleranceIT` and 
its siblings already shut real nodes down in multi-node clusters; adding an 
assertion there (marker present at the coordinator, or the `due to node 
offline` WARN line observed) would have failed on every revision of this PR so 
far. Please add that assertion together with the Issue 1 fix; without it the 
SPI placement cannot be considered verified.
   
   # Compatibility and upgrade
   
   - No wire, checkpoint, savepoint or `TaskExecutionState` format change; 
`FailureClassification` is in-process only. Rolling upgrade is safe: old nodes 
never write a marker, so both consumers take the pre-existing `ERROR` path.
   - Operators with a custom `hazelcast.yaml` keep `TERMINATE` and silently 
keep the old behavior (documented, but see Issue 3 for why the policy change 
should not be in this PR at all).
   - Kubernetes `command` change (`/bin/sh -c "<script>"` to direct script 
exec) and `exec java` in the foreground path are fine: the script already 
carried `#!/bin/bash` and was executed through its shebang before, so no new 
image dependency; the daemon (`-d`) path is untouched.
   
   # Deletion audit
   
   All 42 deleted lines are accounted for: the replaced `errorByPhysicalVertex` 
lines and the old inline `JobException` construction in 
`PhysicalVertex`/`CoordinatorService`, trailing-newline fixes in yaml/md files, 
the `java` to `exec java` line, and the two Kubernetes `command` lines. No 
connector, datasource, registration, config field or fallback was removed.
   
   # Extreme-load assessment
   
   Steady-state cost of the design is one `put` per graceful shutdown and one 
`get` plus one `remove` per member removal on the master, all against a 
TTL-bounded map: negligible. The restore path adds one synchronous `get` per 
vertex whose worker vanished (`PhysicalVertex.java:244-249`), serialized inside 
`initStateFuture`; for jobs with thousands of vertices on one lost worker this 
is thousands of sequential distributed reads during failover. Not blocking, but 
caching the marker per lost address in the restore path would be cheap. The 
blocking risk at scale is Issue 2, not throughput.
   
   # Issue summary
   
   | # | Issue | Severity | Location |
   |---|---|---|---|
   | 1 | Marker write from `GracefulShutdownAwareService.onShutdown()` always 
throws `HazelcastInstanceNotActiveException` (proxy lifecycle check after 
`shuttingDown=true`); WARN downgrade never fires; two new WARN lines per 
graceful shutdown. Confirmed by 0 WARN / 122 ERROR / 153 write failures in the 
last green engine E2E run | P0, blocking | `SeaTunnelServer.java:242,310-343`; 
`Node.java:502-520,633-635`; `AbstractDistributedObject.java:104-115` |
   | 2 | Unguarded `IMap.get` before the failure-propagation loop in 
`failedTaskOnMemberRemoved`; a transient Hazelcast exception now skips failing 
every task on the lost member and can leave the job stuck | P0, blocking | 
`CoordinatorService.java:2179-2184,2197-2226`; `SeaTunnelServer.java:349-357` |
   | 3 | `hazelcast.shutdownhook.policy: GRACEFUL` changes SIGTERM into a 
blocking graceful shutdown (up to 600 s) for every shipped deployment; 
undocumented; unnecessary once Issue 1 is fixed on `SHUTTING_DOWN` | Medium | 
18 config/doc/k8s files; `stop-seatunnel-cluster.sh:51,58` |
   | 4 | Startup-side clear at `STARTING` fails deterministically on 
lite-member workers (`NoDataMemberInClusterException`) and logs a stack trace 
on every worker boot | Medium | `SeaTunnelServer.java:252-256,325-332` |
   | 5 | Cancel path now records an "offline" failure message on `CANCELED` 
states | Low | `PhysicalVertex.java:228-271,457,658-664` |
   | 6 | Javadoc and docs describe an ordering that the source and logs 
disprove | Low | `SeaTunnelServer.java:236-240`; 
`docs/en|zh/engines/zeta/state-storage-and-recovery.md` |
   | 7 | No real-cluster assertion for the marker/WARN path; all new tests are 
mock-based | Low, required with the fix | `ClusterFaultToleranceIT` or a 
dedicated IT |
   
   # CI status on this head
   
   - Fork `Build` run `35119618733` on `fd60dfd596` was still `queued` at 
review time; the apache-side checks on this commit are queued as well. No 
run-verification claim can be made for this head yet.
   - The previous head `e6e243744e` was green (fork run `33966670706`, 83/83 
non-skipped jobs). Green here does not validate this feature: the failing write 
is caught and logged, which is exactly why the E2E logs, not the job status, 
had to be read.
   
   # Merge recommendation
   
   **Blocking, must be fixed before any further review round:**
   1. Issue 1: move the marker write to `stateChanged(SHUTTING_DOWN)` (or prove 
another placement with a real-cluster test), remove the 
`GracefulShutdownAwareService` implementation, and fix the Javadoc/docs (Issue 
6).
   2. Issue 2: make the coordinator's marker read and clear fail-safe, 
mirroring `PhysicalVertex.getGracefulMemberRemovalMarker`.
   3. Issue 7: add a real multi-node assertion for the WARN classification so 
this cannot regress a third time.
   
   **Recommended in the same PR:** drop or split out the 
`hazelcast.shutdownhook.policy` change (Issue 3); move the startup clear to 
`STARTED` (Issue 4); mention the cancel-path message change in the PR body 
(Issue 5).
   
   **Conclusion: implementation deviates from the plan.** The PR body promises 
"only narrows the log level for the known node-offline case" while keeping 
failure handling unchanged. On the real path the log level never narrows, and 
the member-removed failure handling did change, in the wrong direction. I am 
keeping this review as a comment rather than a formal review state because 
GitHub does not allow a formal state on one's own PR. @SEZ9, your original 
Issue 1 was right in substance and is still open on this head: the hook moved, 
but the write still cannot succeed.
   


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