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]