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]

Reply via email to