DanielLeens commented on PR #11856:
URL: https://github.com/apache/seatunnel/pull/11856#issuecomment-5421773819
Thanks for flagging that, @SEZ9 — I checked the raw comment body via the API
and the full text is actually there on GitHub's side (18,939 characters,
nothing truncated in storage), so this looks like a rendering/fetch issue on
your end rather than something lost when I posted it. Re-pasting the two
missing pieces directly so you don't have to chase it further.
## The race (Issue 1, High)
**Root cause:** `PhysicalVertex.updateStateByExecutionService()` writes two
independent atomics as an unsynchronized pair, with mismatched CAS discipline:
```java
// PhysicalVertex.java:546-556
public void updateStateByExecutionService(
TaskExecutionState taskExecutionState, boolean
gracefulMemberRemovalFailure) {
...
errorByPhysicalVertex.compareAndSet(null,
taskExecutionState.getThrowableMsg()); // first-write-wins
gracefulMemberRemovalFailureByPhysicalVertex.set(gracefulMemberRemovalFailure);
// last-write-wins
updateTaskState(taskExecutionState.getExecutionState());
}
```
`errorByPhysicalVertex` is `compareAndSet(null, ...)` — first caller wins
the message, by pre-existing design.
`gracefulMemberRemovalFailureByPhysicalVertex` (new in this PR) is a plain
`.set(...)` — last caller wins the flag. They are not updated as one unit, and
the method has no `synchronized` guard (unlike
`updateTaskState()`/`stateProcess()`, which are both `synchronized`).
**Concrete interleaving:**
```
Thread A (Hazelcast operation-executor thread — a real task failure reported
over RPC):
NotifyTaskStatusOperation
-> JobMaster.updateTaskExecutionState(taskExecutionState)
[JobMaster.java:1253/1267]
-> physicalVertex.updateStateByExecutionService(taskExecutionState)
// 1-arg overload
-> updateStateByExecutionService(taskExecutionState, false)
[PhysicalVertex.java:539]
errorByPhysicalVertex.compareAndSet(null, "<real exception
message>") // succeeds
gracefulMemberRemovalFailureByPhysicalVertex.set(false)
Thread B (Hazelcast membership-event thread — same vertex, classified
node-offline because
this same worker is also mid-graceful-shutdown around the same time):
CoordinatorService.failedTaskOnMemberRemoved(event)
[CoordinatorService.java:2030]
-> makeTasksFailed(..., gracefulMemberRemoval=true)
[CoordinatorService.java:2057]
(guard: executionState == RUNNING — still true if Thread A hasn't
transitioned yet)
-> physicalVertex.updateStateByExecutionService(
buildMemberRemovedFailureState(...), true)
[CoordinatorService.java:2071]
errorByPhysicalVertex.compareAndSet(null, "...") // FAILS
silently, A already won
gracefulMemberRemovalFailureByPhysicalVertex.set(true) //
overwrites A's `false`
```
If Thread B's `.set(true)` lands after Thread A's `.set(false)` — entirely
possible, since `NotifyTaskStatusOperation` (operation-executor pool) and
`MembershipAwareService.memberRemoved` (cluster/membership-event thread) are
different Hazelcast thread pools with no ordering guarantee between them — then
whichever thread wins the race into `synchronized updateTaskState()` and
actually flips `RUNNING -> FAILED` calls `stateProcess()`, which reads:
- `errorByPhysicalVertex.get()` = Thread A's real exception message (correct)
- `gracefulMemberRemovalFailureByPhysicalVertex.get()` = `true`, from Thread
B
Net effect: a genuine task failure, carrying its own genuine error message,
gets logged at `log.warn` instead of `log.error`
(`PhysicalVertex.java:632-641`). The message text is still the real exception,
so a human reading the raw line isn't misled about *what* failed — but the
severity, which is exactly the signal this PR exists to make trustworthy, is
wrong. The losing thread hits the pre-existing `current.isEndState()` guard and
returns silently, so there's no double-log to notice the discrepancy from.
**Fix options:**
- Option A (more robust): collapse both fields into one
`AtomicReference<FailureClassification>` (an immutable `{message,
gracefulFlag}` holder) updated via a single `compareAndSet(null, ...)`, so
message and flag always travel together under first-write-wins.
- Option B (minimal diff): only write the flag when the paired message-write
actually wins:
```java
if (errorByPhysicalVertex.compareAndSet(null,
taskExecutionState.getThrowableMsg())) {
gracefulMemberRemovalFailureByPhysicalVertex.set(gracefulMemberRemovalFailure);
}
```
One-line guard, restores "the flag belongs to whichever call's message
actually stuck."
I lean toward Option A since this is state future maintainers will likely
extend, but Option B closes the actual bug on its own.
## CI status at review time
`Build` was failing on fork run `32628391769` (head `06344bd0`). Pulled the
job log directly: every module through `seatunnel-engine-server` (the module
this PR touches) builds and unit-tests green. The only failing job is
`rocketmq-connector-it (11, ubuntu-latest)`, one failed test —
`RocketMqIT.testSourceRocketMqRestore`, `Expected 20 '_initial_' messages, got:
25` — preceded in the log by `RemotingSendRequestException: send request to
<172.18.0.2:9876> failed` against the RocketMQ test container. That reads as a
message-count/duplicate-delivery flake in an unrelated connector's E2E suite,
not something a `seatunnel-engine`-only diff could cause. Still needs to go
green (or be explicitly waived as known flake) before merge, per the usual bar.
Full text is also on the original review if the render issue clears up:
https://github.com/apache/seatunnel/pull/11856#pullrequestreview-5009299593
Once the pairing fix (Option A or B) is in, I'll re-verify against the new
head.
--
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]