DanielLeens commented on PR #11856:
URL: https://github.com/apache/seatunnel/pull/11856#issuecomment-5711741991
[REVIEW] Maintainer re-review of head `b91afe70e1`, two commits after the
head I reviewed yesterday (`fd60dfd596`): `f3308d8168` "Correct graceful member
removal lifecycle" and `b91afe70e1` "Remove no-op files from PR diff". I
re-read the full diff against `dev` (now 14 files, +1245/-34) and the
incremental diff, and I read the CI logs of the run on this head rather than
its status only.
**Verdict: not ready to merge. The design fixes for both P0s from the last
round are correct at source level, but this head does not compile its test
sources, so CI is red on 65 jobs and the fix has no runtime verification at
all.** The Validation section of the PR body claims "real two-node engine E2E
coverage proving that `SHUTTING_DOWN` publishes a marker". On this head that
test has never executed.
# 1. Blocking: the head does not compile (N1)
```text
[ERROR] COMPILATION ERROR :
[ERROR]
.../seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/SeaTunnelServerShutdownTest.java:[72,27]
error: cannot find symbol
symbol: variable TimeUnit
location: class SeaTunnelServerShutdownTest
[ERROR] Failed to execute goal ...maven-compiler-plugin:3.10.1:testCompile
(default-testCompile) on project seatunnel-engine-server
```
- Cause: `f3308d8168` removed `import java.util.concurrent.TimeUnit;` from
`SeaTunnelServerShutdownTest` while line 72 still asserts
`eq(TimeUnit.MILLISECONDS)` on the marker `put`. The fix is to restore that one
import.
- Scope: fork run `35175217738` on this head has 65 failed jobs. I opened
five of them (`unit-test (8)`, `engine-v2-it (8)`, `kafka-connector-it (8)`,
`jdbc-connectors-it-part-1 (8)`, `transform-v2-it-part-1 (11)`); all five stop
at this exact error. The 9 green jobs are the ones that never compile Java test
sources (license, dead links, code style, Helm, website, UI, dependency
licenses, sanity, changes). Two `oracle-cdc` jobs were still running and the
apache-side `Build` check was pending when I looked.
- Consequence:
`ClusterFaultToleranceIT.testGracefulShutdownPublishesMemberRemovalMarker`, the
new unit tests, and every existing fault-tolerance IT did not run. The local
validation listed in the PR body (scoped spotless plus `git diff --check`)
cannot catch a missing import, so please do not describe the E2E as "proving"
anything until a CI run has executed it.
# 2. Status of yesterday's findings on this head
| # | Finding from the last round | Status on `b91afe70e1` | Evidence |
|---|---|---|---|
| 1 (P0) | Marker write from `GracefulShutdownAwareService.onShutdown()` can
never succeed | **Fixed in source, unverified at runtime.**
`GracefulShutdownAwareService` is gone; `stateChanged(SHUTTING_DOWN)` now
writes the marker | `SeaTunnelServer.java:239-245`, `:300-326`; listener
registered at `:154`. `LifecycleServiceImpl.shutdown()` fires `SHUTTING_DOWN`
at `:94` before `node.shutdown()` at `:101`, so `Node.isRunning()` is still
`true` there |
| 2 (P0) | Unguarded `IMap.get` ahead of the failure-propagation loop |
**Fixed.** Read and clear are wrapped and fall back to the unproven/`ERROR`
path | `CoordinatorService.java:2180-2195` (read), `:2202-2215` (clear); tests
`shouldTreatMarkerReadFailureAsUnproven`, `shouldAbsorbMarkerClearFailure` |
| 3 | `hazelcast.shutdownhook.policy: GRACEFUL` changed SIGTERM semantics
everywhere | **Resolved.** All 17 config/doc/k8s-conf files are byte-identical
to the merge base again; SIGTERM behaves exactly as on `dev` | `git diff
--quiet <merge-base> HEAD -- config deploy/.../conf ...` reports identical |
| 4 | Startup clear at `STARTING` failed on every lite-member worker boot |
**Fixed.** Clear moved to `STARTED`, which `HazelcastInstanceFactory.java:241`
fires only after the join | `SeaTunnelServer.java:242-244` |
| 5 | Cancel path now records an "offline" message | **Acknowledged** in the
PR body under "Behavior and compatibility" | PR description |
| 6 | Javadoc/docs described the wrong ordering | **Fixed.** Javadoc and
both `state-storage-and-recovery.md` pages now describe `SHUTTING_DOWN` and the
default `TERMINATE` policy | docs diff |
| 7 | No real-cluster test | **Added, never executed** (see N1) |
`ClusterFaultToleranceIT.testGracefulShutdownPublishesMemberRemovalMarker` |
One improvement beyond what I asked for: the coordinator clear is now
value-conditional, `remove(lostAddress, markedAt)` at
`CoordinatorService.java:2207`, so it cannot delete a newer marker written by a
restarted member at the same address. Good change.
I also checked the claim in the new docs and PR body that only intentional
exits publish a marker. It holds in Hazelcast 5.1: the callers that go through
`LifecycleService` are `NodeShutdownHookThread` (`Node.java:766,769`),
`HazelcastInstanceImpl.shutdown()` (`:346`), the factory's
`shutdownAll`/`terminateAll`, and the operator cluster-shutdown API
(`ClusterServiceImpl.java:959,986`, `ShutdownNodeOp.java:44`). Fault paths
bypass it and call `node.shutdown(true)` directly, so they emit no
`SHUTTING_DOWN`: `OutOfMemoryHandlerHelper.tryShutdown` (`:62`, which is where
SeaTunnel's `ExceptionUtil` routes an `OutOfMemoryError` via
`OutOfMemoryErrorDispatcher`) and `ClusterServiceImpl.java:400`. A JVM OOM
therefore stays at `ERROR`, as it should.
# 3. New non-blocking findings
**N2 (Medium, recommended): both lifecycle hooks now do synchronous
distributed calls on lifecycle threads.**
- `stateChanged(SHUTTING_DOWN)` runs a blocking `IMap.put` inside
`LifecycleServiceImpl.shutdown()`'s `synchronized (lifecycleLock)` block
(`LifecycleServiceImpl.java:93-94`), on the caller's thread. Under the default
`TERMINATE` policy that caller is the JVM shutdown-hook thread, in front of a
`terminate()` that is expected to be immediate. `stateChanged(STARTED)` runs a
blocking `IMap.remove` on the thread that constructs the instance.
- Both are bounded, and both fail safe. With the shipped
`hazelcast.invocation.max.retry.count: 20`, retryable errors give up after
roughly 3.5 s (`Invocation.handleRetry`, `Invocation.java:706-712`). The long
tail is a partition owner that is alive but unresponsive: that waits for the
operation call timeout, which is the same ~120 s stall this PR already hit once
at `45e981d3`.
- This is an inference from the source, not something I observed. Operators
restart nodes precisely when the cluster is degraded, so I would not let a
best-effort log hint hold the shutdown hook for two minutes. Suggestion:
`putAsync(...)` / `removeAsync(...)` with a short bounded wait of a few
seconds, keeping the existing catch-and-WARN fallback. On a full-cluster stop
some nodes will also log `Failed to mark graceful member removal marker` with a
stack trace because their peers are already gone; consider dropping the stack
trace for that known case.
**N3 (Low): the new E2E proves publication, not classification.** It shuts
down `node1` with no job running and asserts that `node2` sees the marker
through a real proxy. That is the right proof for the defect that was actually
broken, and using an entry listener instead of polling `get` is correct,
because the surviving coordinator clears the marker right after
`memberRemoved`. The end-to-end evidence that the log level really narrows
should come from the next engine job log, where the existing worker-down ITs
already produce these lines. Acceptance criteria I will check myself on the
next green run, against yesterday's baseline (153 write failures, 0 WARN, 122
ERROR in `engine-v2-it (8)`):
- `Failed to mark graceful member removal marker` followed by
`HazelcastInstanceNotActiveException` is gone for `shutdown()` paths;
- `end with state FAILED due to node offline` (WARN) appears for the
worker-down tests;
- `Failed to clear graceful member removal marker` with
`NoDataMemberInClusterException` is gone.
**N4 (Low, docs): marker TTL versus failure detection.** The TTL is 300 s
(`Constant.java:72`) and the shipped configs detect a silent member after
`hazelcast.max.no.heartbeat.seconds: 180`. An operator who raises that timeout
above 300 s will see heartbeat-detected departures classified as `ERROR`
because the marker expires first. It fails safe, so one sentence in
`state-storage-and-recovery.md` is enough.
# 4. @SEZ9, on your comment from 02:55Z
Your comment was written against `fd60dfd596` and lists F1 plus F2 through
F8. F2 through F8 were closed in earlier rounds; I re-verified each on this
head so the thread has one clear status:
- **F3/F6 payload:** `buildMemberRemovedFailureState` wraps the message in
`JobException` on both paths (`CoordinatorService.java:2144-2149`), pinned by
`shouldKeepThrowablePayloadForMemberRemovedFailureState`. No plain-string
branch exists any more.
- **F2/F4 marker lifecycle:** the read is a non-destructive `get` (`:2185`);
the clear runs after the propagation loop, only when no master-switch restore
is in flight, and is value-conditional (`:2207`); the write uses Hazelcast's
native TTL overload (`SeaTunnelServer.java:309-313`). The timestamp check is a
secondary guard with a ±5 min skew tolerance that fails toward `ERROR`.
- **F5/F8 message matching:** there is no regex left in `PhysicalVertex` (no
`Pattern`, no `.matches(`). Classification is a typed boolean carried through
`updateStateByExecutionService(state, graceful)` and the immutable
`FailureClassification` holder, and the offline message has a single producer,
`CoordinatorService.buildMemberRemovedOfflineMessage`.
- **F7 docs:** present in
`docs/en|zh/engines/zeta/state-storage-and-recovery.md`.
- **F1:** your diagnosis was right. It is fixed in source on this head and
unverified at runtime because of N1. The log excerpt you asked for is exactly
the N3 acceptance list.
# 5. Compatibility, deletions, load
- Unchanged from the last round: no wire, checkpoint, savepoint or
`TaskExecutionState` format change; old members never write a marker and both
consumers then take the `dev` `ERROR` path, so rolling upgrade is safe. With
Issue 3 reverted there is no longer any change to shutdown semantics or shipped
configuration.
- Deletion audit: the 34 deleted lines are the replaced
`errorByPhysicalVertex` handling and inline `JobException` construction, `java`
to `exec java`, and the two Kubernetes `command` lines. No connector,
registration, config field or fallback was removed.
- Load: one `put` per intentional shutdown, one `get` plus one conditional
`remove` per member removal, all on a TTL-bounded map. The per-vertex marker
read on the master-failover restore path is unchanged and still not blocking.
# 6. Merge recommendation
**Blocking:**
1. N1: restore the `TimeUnit` import and get a complete CI run on the fixed
head.
2. After that run, the N3 evidence must be present in the `engine-v2-it`
log. Green status alone is not accepted for this PR; that is what let the inert
implementation through last time.
**Recommended in the same PR:** N2 (bounded async marker calls).
**Optional:** N4.
**Conclusion: implementation deviates from the plan.** At source level the
change now matches what the PR describes, and both P0s are addressed the right
way. The deviation is the Validation claim: the head does not compile, so none
of the tests it cites have run. Posted as a comment because GitHub does not
allow a formal review state on one's own PR.
--
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]