DanielLeens commented on PR #11856:
URL: https://github.com/apache/seatunnel/pull/11856#issuecomment-5379365917
Since this is my own PR, GitHub won't let me submit a formal review state on
it, so this is posted as a plain comment again. This is round 7 on this PR (my
sixth pass plus @SEZ9's independent review). The code itself is unchanged since
my last comment (head is still `f953140ec`), so I'm not re-deriving the whole
design from scratch again — I independently re-verified the current source
against every load-bearing claim from round 6 and @SEZ9's review (all confirmed
accurate), and I re-checked CI live rather than reusing round 6's "still
queued" status. That live check changes the conclusion: **CI has since
completed and the current head does not compile.** This is a hard blocker that
supersedes my own round-6 "Ready to merge after fixes."
# What Problem Does This PR Solve?
When a worker leaves the Hazelcast cluster during a routine,
operator-initiated scale-down, `CoordinatorService` fails its deployed task
groups and `PhysicalVertex.stateProcess()` always logged that at `ERROR`,
making expected scale-down noise indistinguishable from a real engine fault.
After six rounds of iteration (fixing a message-construction bug that made the
original classification a no-op, narrowing "graceful" from "any member removal"
to a self-reported per-address marker, gating the marker on Hazelcast's
`terminate` flag, and — following @SEZ9's finding — relocating the marker
publication to the correct pre-shutdown SPI hook), the mechanism now publishes
a best-effort "I'm leaving on purpose" marker into a new
`engine_gracefulMemberRemoval` IMap before a graceful shutdown, consumes it
when classifying the resulting task failure, and downgrades only the exact
node-offline message to `WARN` when the marker proves the departure was
graceful. Any unproven/abrupt departure
(crash, OOM, forced termination, partition) still logs `ERROR`, unchanged
from today.
# 1. Code Change Review
## 1.1 Core Logic Analysis
I independently re-verified the current diff against the tree directly (not
by re-reading the prior rounds' write-ups) and confirm every load-bearing claim
from round 6 and @SEZ9's review holds:
- `SeaTunnelServer.java`: `SeaTunnelServer` now additionally implements
`GracefulShutdownAwareService`; `onShutdown(long, TimeUnit)` calls
`markLocalGracefulMemberRemoval()`, and the marker-write call is removed
entirely from `ManagedService.shutdown(boolean terminate)` (confirmed by
reading the full `shutdown()` body — every remaining line there operates on
separately-mocked, null-guarded fields, none of them `nodeEngine`). `init()`
clears this node's own stale marker on startup.
- `CoordinatorService.java`: adds `gracefulMemberRemovalIMap`, wired in
`initCoordinatorService()`;
`buildMemberRemovedOfflineMessage`/`buildMemberRemovedFailureState`/`isGracefulMemberRemovalMarkerValid`/`consumeGracefulMemberRemovalMarker`
are exactly as described — the graceful branch builds `TaskExecutionState` via
the `String` constructor (plain message, regex-matchable), the non-graceful
branch still wraps in `new JobException(...)` (full stack trace, never matches,
`ERROR` unchanged). `failedTaskOnMemberRemoved` calls
`consumeGracefulMemberRemovalMarker` once, up front, and threads the resulting
boolean through both `makeTasksFailed` call sites.
- `PhysicalVertex.java`: `DEPLOYED_NODE_OFFLINE_ERROR_PATTERN` and
`isDeployedNodeOfflineFailure()` are exactly as described; `stateProcess()`'s
`FAILED` case now branches `log.warn` vs `log.error` on that predicate.
Confirmed the message-template is genuinely defined in two disconnected places
with no shared constant (`CoordinatorService.java`'s `%s`/`%s` format string
vs. `PhysicalVertex.java`'s hand-written regex) — @SEZ9's Issue 8 is real and I
could find no test that feeds one into the other.
- `Constant.java`: adds `IMAP_GRACEFUL_MEMBER_REMOVAL` and a 5-minute
`GRACEFUL_MEMBER_REMOVAL_MARK_TTL_MILLIS`, both with adequate Javadoc this
round (closing part of round-6's Issue 3 from earlier rounds).
**Runtime path** (unchanged from round 6's trace, re-confirmed against
current source):
```text
Graceful path:
operator-initiated shutdown -> Hazelcast Node.shutdown(false)
-> callGracefulShutdownAwareServices() -> SeaTunnelServer.onShutdown()
[SeaTunnelServer.java:233-237]
markLocalGracefulMemberRemoval() ->
gracefulMemberRemovalIMap.put(thisAddress, now)
-> (network) coordinator observes memberRemoved
-> CoordinatorService.failedTaskOnMemberRemoved(event)
[CoordinatorService.java:2039]
gracefulMemberRemoval =
consumeGracefulMemberRemovalMarker(lostAddress) [:2008-2013]
-> makeTasksFailed(..., gracefulMemberRemoval)
-> buildMemberRemovedFailureState(...): plain-String
TaskExecutionState when graceful
-> PhysicalVertex.stateProcess() case FAILED:
isDeployedNodeOfflineFailure(errorMsg) matches -> log.warn
Abrupt path (kill -9 / OOM / partition / forced terminate):
no marker ever written for that address ->
consumeGracefulMemberRemovalMarker returns false
-> buildMemberRemovedFailureState(..., false) ->
JobException(offlineMessage) -> stack trace
-> isDeployedNodeOfflineFailure() never matches a multi-line stack trace
-> log.error (unchanged from dev)
```
**What I found that neither round 6 nor @SEZ9 caught: the current head does
not compile.** I pulled the live CI run for the current head SHA
(`f953140ecf6ee1e96544679c0ed28d91e51d94b0`, `DanielLeens/seatunnel` fork run
`32484015488`, completed 2026-08-22, `conclusion: failure`) rather than
trusting round 6's "still queued at review time" note. All four `unit-test`
matrix legs (JDK 8/11 x ubuntu/windows) fail, and I pulled the raw job log to
find the actual cause rather than assuming it was another Windows/CDC-style
environmental flake:
```
[ERROR] COMPILATION ERROR :
[ERROR]
.../seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/SeaTunnelServerShutdownTest.java:[55,16]
error: no suitable method found for thenReturn(IMap<Address,Long>)
[ERROR] method OngoingStubbing.thenReturn(IMap<Object,Object>) is not
applicable
[ERROR] (argument mismatch; IMap<Address,Long> cannot be converted to
IMap<Object,Object>)
```
This is a `mvn test-compile` failure in `seatunnel-engine-server`, at:
```java
IMap<Address, Long> gracefulMemberRemovalIMap = mock(IMap.class); //
raw-typed mock, unchecked cast
...
when(hazelcastInstance.getMap(Constant.IMAP_GRACEFUL_MEMBER_REMOVAL))
.thenReturn(gracefulMemberRemovalIMap); // javac
can't reconcile the generics here
```
`mock(IMap.class)` produces a raw `IMap`; assigning it to the parameterized
local variable is an unchecked conversion, and when it's then passed into
`.thenReturn(...)` on the `OngoingStubbing<IMap<Object,Object>>` inferred from
`HazelcastInstance.getMap(String)`'s unbound type parameters, javac 8/11 both
fail to unify `IMap<Address,Long>` with `IMap<Object,Object>`. I confirmed this
same single root cause cascades into the entire CI matrix, not just
`unit-test`: I independently pulled the log for `all-connectors-it-1 (8,
ubuntu-latest)` — a job with zero topical relation to this PR — and it shows
the identical `COMPILATION ERROR` at the identical line, because that job's
Maven invocation also builds `seatunnel-engine-server` with `-am` in the shared
reactor. That is why essentially the entire 60+-job CI matrix is red on this
head, not because of 60 unrelated flakes.
**Key findings:**
- The design and runtime-path logic (the actual subject of six rounds of
review) remain sound and unchanged — the compile break is confined to one test
file and does not reflect a defect in the production code path.
- This is a real, reproducible compile failure on this exact head, not a
flake — it is 100% deterministic (a `javac` generics-inference limitation, not
a race or environment issue), and it fully explains the "mass CI failure"
pattern rather than requiring 60 separate unrelated-flake diagnoses.
- The fix is small and well-understood: either give `getMap` an explicit
type witness (`hazelcastInstance.<Address, Long>getMap(...)`), or stub via
`doReturn(gracefulMemberRemovalIMap).when(hazelcastInstance).getMap(...)`
(which sidesteps `OngoingStubbing`'s generic-return-type inference entirely and
is the more common idiom elsewhere in this codebase for exactly this class of
Mockito+generics friction).
## 1.2 Compatibility Impact
Fully compatible. No config `Option`, public API, SPI contract break
(implementing an additional Hazelcast-defined interface is additive), or
checkpoint/wire format touched. Rolling-upgrade safety is preserved: an
old-version member never implements `GracefulShutdownAwareService`, so its
departure is never marked graceful and falls back to the pre-existing
always-`ERROR` behavior identical to current `dev`.
## 1.3 Performance / Side-Effect Analysis
Negligible — one `IMap.put()` on graceful shutdown (now dispatched from
Hazelcast's own bounded `graceful-shutdown` executor rather than inline in
`ManagedService.shutdown()`), one `IMap.remove()` per `memberRemoved` event on
the already-slow-path of task failure handling. No new locks, no hot-path
allocation.
## 1.4 Error Handling and Logging
**Issue 1 — Current head does not compile; CI is red on a genuine compile
error, not a flake**
- Location:
`seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/SeaTunnelServerShutdownTest.java:55`
- Problem: `mock(IMap.class)` (raw type) assigned to `IMap<Address, Long>`
and then passed to `.thenReturn(...)` fails javac's generic-return-type
inference against `HazelcastInstance.getMap(String)`'s unbound `<K, V>`
parameters, on both JDK 8 and JDK 11.
- Potential risk: this is not cosmetic — it blocks every CI job in the
matrix (confirmed identical failure in an unrelated `all-connectors-it-1` job
via the shared `-am` reactor build), and it must be treated as a
release-blocking correctness gate regardless of how sound the surrounding
design is.
- Best improvement:
`Mockito.doReturn(gracefulMemberRemovalIMap).when(hazelcastInstance).getMap(Constant.IMAP_GRACEFUL_MEMBER_REMOVAL);`
(avoids `OngoingStubbing`'s generic return-type inference entirely), or add an
explicit type witness: `hazelcastInstance.<Address, Long>getMap(...)`.
- Severity: Critical
- Raised by another reviewer: No — CI was still queued at the time of both
my round-6 review and @SEZ9's review; neither of us could have caught this from
source reading alone since the code compiles under simple visual inspection
(the type mismatch is a javac inference artifact, not a type error a human eye
reliably catches). This is a genuine miss on my part in round 6 for asserting
"no blockers remaining... contingent on CI" without following up once CI
actually completed. I'm not going to let that stand uncorrected.
**Issue 2 (carried over from round 6, unaddressed) — `onShutdown()`'s
comment misstates the mechanism it relies on**
- Location: `SeaTunnelServer.java:233-235`
- The comment says "before the node becomes PASSIVE"; per Hazelcast 5.1's
`Node#shutdown(boolean)`, `NodeState` is set to `PASSIVE` unconditionally
before `callGracefulShutdownAwareServices()` runs, so the node is already
`PASSIVE` at that point — the actual reason the placement is safe is that it
runs before `NodeEngineImpl.shutdown()` tears down
`proxyService`/`operationService` (which back `getMap()`). I verified this
against the actual Hazelcast 5.1 source in round 6 and it's still uncorrected
in the current diff.
- Risk: low on its own; a future maintainer reasoning from the wrong premise
could relocate this call to somewhere that again races proxy-service teardown.
- Best improvement: reword to describe the real mechanism (service-shutdown
ordering, not `NodeState`).
- Severity: Low
- Raised by another reviewer: Partially (@SEZ9 raised the underlying
placement concern; the comment-accuracy point is mine).
**Issue 3 (carried over, unaddressed) — marker consumption is destructive
with no Hazelcast-native TTL**
- Location: `CoordinatorService.java` (`consumeGracefulMemberRemovalMarker`,
`remove()` before any task group is failed), `SeaTunnelServer.java` (`put()`
with no TTL overload).
- A master failover between `remove()` and the end of the propagation loop,
or any departure whose marker is never read back, leaves an unbounded,
un-evicted entry.
- Best improvement: use the TTL-bearing `put(key, value, ttl, unit)`
overload; consider non-destructive read until the propagation loop completes.
- Severity: Medium
- Raised by another reviewer: Yes (@SEZ9, Issues 2 and 4).
**Issue 4 (carried over, unaddressed) — graceful branch swaps the
`TaskExecutionState` payload from `JobException` to a bare `String`**
- Location: `CoordinatorService.buildMemberRemovedFailureState`.
- Downstream consumers of `throwableMsg` (job history, REST error fields)
see a differently-shaped payload depending on which branch fired.
- Best improvement: keep `JobException` on both branches; carry the
classification as an explicit flag instead of a payload-shape difference.
- Severity: Medium
- Raised by another reviewer: Yes (@SEZ9, Issues 3 and 6).
**Issue 5 (carried over, unaddressed) — offline-message template duplicated
with no shared constant or round-trip test**
- Location: `CoordinatorService.buildMemberRemovedOfflineMessage` (format
string) vs. `PhysicalVertex.DEPLOYED_NODE_OFFLINE_ERROR_PATTERN` (regex).
- I re-confirmed directly against the current tests: `PhysicalVertexTest`
uses a hand-written literal, `CoordinatorServiceMemberRemovedTest` never feeds
`buildMemberRemovedOfflineMessage`'s actual output into
`isDeployedNodeOfflineFailure`. No test today would catch drift between the two.
- Best improvement: derive both from one shared constant, or add a
round-trip test.
- Severity: Medium
- Raised by another reviewer: Yes (@SEZ9, Issue 8). This is the cheapest,
highest-leverage of the non-blocking follow-ups.
**Issue 6 (carried over, unaddressed) — classification by free-text regex
crosses a trust boundary**
- Location: `PhysicalVertex.isDeployedNodeOfflineFailure`, fed by
`errorByPhysicalVertex`, which can also be populated from
worker-reported/connector-originated messages.
- A connector exception whose message happens to match the exact template
would be silently downgraded to `WARN`.
- Best improvement: carry the classification as a typed flag on
`TaskExecutionState`, not message text.
- Severity: Medium
- Raised by another reviewer: Yes (@SEZ9, Issue 5).
**Issue 7 (carried over, unaddressed) — no docs/en / docs/zh update for the
new IMap / WARN-vs-ERROR classification**
- Severity: Low
- Raised by another reviewer: Yes (@SEZ9, Issue 7).
# 2. Code Quality Assessment
## 2.1 Coding Standards
ASF headers present, Javadoc present on all new non-trivial fields/methods
(a real improvement over round 2's gap), no wildcard imports, no
`System.out.println`. The one standards issue is Issue 1 above — this specific
test does not compile.
## 2.2 Test Coverage and Test Stability
The classification logic (`isDeployedNodeOfflineFailure`), the
message-construction branching (`buildMemberRemovedFailureState`), and the
marker-write/no-write branching in `SeaTunnelServerShutdownTest` are all
exercised as plain, synchronous, mock-based unit tests — no `Thread.sleep`, no
shared static state, no timing dependence, no containers/network I/O. Absent
the compile error, this would rate **Stable**. Because of Issue 1, **the actual
current rating must be: the test suite does not build, and therefore provides
zero coverage on this head** — a compile failure is a stronger negative signal
than any flaky-test pattern, since it means none of these otherwise
well-designed tests actually ran. I'm not going to soften this by rating the
*design* of the tests separately from the fact that they don't compile; the
mandatory rating for this PR as it stands is **High risk**, purely on the
compile break, with a clear and narrow one-line fix.
The remaining, pre-existing coverage gap (unchanged from round 6): no
integration-level test exercises the real distributed `IMap` write/read round
trip through an actual multi-node Hazelcast cluster — both my round-2 review
and @SEZ9 flagged this; still open, non-blocking.
## 2.3 Documentation Updates
Not applicable for this specific commit's production code; see Issue 7 for
the pre-existing, PR-wide documentation gap (operator-visible WARN/ERROR
classification change, new IMap).
# 3. Architectural Soundness
## 3.1 Elegance of the Solution (Precise fix / Temporary workaround /
Long-term solution)
The design itself is a **precise fix** for the stated problem, verified
independently across six rounds including against actual upstream Hazelcast 5.1
source for the SPI-ordering question @SEZ9 raised. It remains fundamentally a
message-text classification at the `PhysicalVertex` boundary (Issues 4-6),
which is a reasonable minimal-diff tradeoff for this PR's scope, not a design
defect. The one thing keeping this out of "long-term solution" territory today
is purely the build break in Issue 1, which is mechanical, not conceptual.
## 3.2 Maintainability
Good and improving round over round, aside from the still-open
comment-accuracy nit (Issue 2) and the message-duplication risk (Issue 5) that
a future edit could silently regress.
## 3.3 Extensibility
Unchanged from round 6 — the self-reported marker pattern generalizes to
future "was this expected?" classifications if needed.
## 3.4 Historical-Version Compatibility
No impact — no serialized state, checkpoint, or wire format touched; old,
marker-unaware members simply never get a marker, unchanged always-`ERROR`
fallback.
# 4. Issue Summary
| # | Issue | Location | Severity | Raised by another reviewer |
|---|---|---|---|---|
| 1 | Current head fails `test-compile` in `seatunnel-engine-server`
(Mockito/javac generics inference), cascading a red CI matrix across ~60
unrelated jobs via the shared `-am` reactor build |
SeaTunnelServerShutdownTest.java:55 | Critical | No |
| 2 | `onShutdown()`'s comment misstates the real mechanism (claims "before
PASSIVE"; actually before proxy/operation-service teardown, while already
PASSIVE) | SeaTunnelServer.java:233-235 | Low | Partially (@SEZ9) |
| 3 | Marker consumption destructive, no Hazelcast-native TTL;
failover/never-consumed entries can leak or misclassify |
CoordinatorService.java, SeaTunnelServer.java | Medium | Yes (@SEZ9) |
| 4 | Graceful branch swaps `JobException` payload for a bare `String`,
inconsistent with the sibling branch | CoordinatorService.java | Medium | Yes
(@SEZ9) |
| 5 | Offline-message template duplicated (format string vs. regex), no
shared constant or round-trip test | CoordinatorService.java,
PhysicalVertex.java | Medium | Yes (@SEZ9) |
| 6 | Classification by free-text regex match crosses a trust boundary |
PhysicalVertex.java | Medium | Yes (@SEZ9) |
| 7 | No docs/en / docs/zh update for the new IMap / WARN-vs-ERROR
classification | N/A | Low | Yes (@SEZ9) |
# 5. Merge Recommendation
### Conclusion: Not recommended for merge
**Blockers — must be fixed:**
1. Issue 1 (Critical) — fix the Mockito/generics compile break in
`SeaTunnelServerShutdownTest.java:55` (`doReturn(...).when(...)` or an explicit
type witness on `getMap`), then get a fully green CI run on the resulting head
before merge. This is the sole reason I'm not endorsing my own round-6 "Ready
to merge after fixes" today — that conclusion was reached while CI was still
queued, and the completed run tells a different story than I assumed it would.
**Recommended fixes — non-blocking, all carried over and still valid:**
- Issue 5 — add a round-trip test (or shared constant) tying the message
format string to the regex; cheapest, highest-leverage of the remaining items.
- Issue 3 — native-TTL `IMap.put` overload, non-destructive-until-processed
marker consumption.
- Issue 4 / Issue 6 — carry the graceful/non-graceful classification as an
explicit typed flag rather than a payload-shape or free-text difference; would
close both at once.
- Issue 2 — correct the `onShutdown()` comment.
- Issue 7 — a short docs note on the WARN/ERROR distinction for operators
who alert on Zeta engine logs.
**Overall assessment:** the design work across six rounds — including a
second independent reviewer's real catch and a verified-correct fix for it —
has been genuinely solid, and I stand by everything I said about the runtime
logic in round 6: it fails safe (every ambiguous case defaults to the
pre-existing `ERROR` behavior), and the mechanism itself is sound. But "ready
to merge after fixes, contingent on CI" was not actually satisfied by this head
— CI finished, and it's red for a real, deterministic reason. I should have
circled back and confirmed the run before writing that conclusion instead of
asking a maintainer to check it independently; I'm correcting that now rather
than letting a stale "no blockers" stand on the PR. Once Issue 1 is fixed and
CI is actually green on the resulting commit, I'd be comfortable with a
maintainer merge decision — the remaining Medium/Low items are legitimate
hardening work, not correctness gates.
--
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]