DanielLeens commented on PR #12377: URL: https://github.com/apache/seatunnel/pull/12377#issuecomment-5726070148
Thanks for digging into this so carefully, @SEZ9 — you're right that there's a real gap here that my own review didn't pin down precisely enough. I re-checked the current head (`24b817941c53`) rather than just re-reading my earlier notes. **Where your trace is exactly right:** for a vertex whose worker crashes *before* `PhysicalVertex.cancel()` has been invoked on it (still `RUNNING`), `CoordinatorService.failedTaskOnMemberRemoved()` -> `makeTasksFailed()` -> `updateStateByExecutionService(FAILED)` can win the race and complete that vertex's future as `FAILED` before the main cancel loop even reaches it — this is exactly the "vertex that hasn't been told to cancel yet... can win the race outright" case I flagged as a residual risk in my own review but didn't push hard enough on. Once that happens, `addPhysicalVertexCallBack` (SubPlan.java:216-219) increments `failedTaskNum` and calls `updatePipelineState(FAILING)` before the pipeline is ever observed as `CANCELING` by `determinePipelineEndState()`, so for that ordering the pipeline does not end `CANCELED`. That's a genuine, unaddressed gap, not a hypothetical. **Where I don't think "unreachable in the real callback path" holds, though:** for a vertex on which `cancel()` has *already* been called (the actual #12353 scenario, and what the regression test exercises), `PhysicalVertex.cancel()` -> `stateProcess()` -> `noticeTaskExecutionServiceCancel()` (PhysicalVertex.java:409-463) are all `synchronized` on that vertex's own monitor, and the method holds that lock through the entire retry-and-sleep loop until it either gets an ack or detects the member is gone and self-resolves to `CANCELED` at line 463. A concurrent `updateStateByExecutionService(FAILED)` for that same vertex (from `makeTasksFailed`) blocks on the same monitor, and once it finally acquires it, hits the terminal-state guard in `updateTaskState` (PhysicalVertex.java:376-379, `current.isEndState()`) and is discarded — the vertex is already `CANCELED` by then. That's the ordering `SplitClusterFaultToleranceIT#testStreamJobCancelResolvesWhenWorkerCrashesBeforeCancelAck` cover s, which is why I pulled the fork's CI logs myself and saw it pass 10/10 on both JDK 8 and JDK 11 rather than taking the PR description's word for it. So, reconciling both: the fix closes the race for "worker crashes on a vertex that's already mid-cancel" (what's tested today), but not for "a sibling vertex on the same crashed worker that the cancel loop hasn't reached yet" (what you traced). Both are real production timings for a multi-task pipeline that loses a worker running more than one of its tasks. Given that, I think your suggested fix direction is the right one to close the remaining gap. Between your two options, skipping the `updatePipelineState(FAILING)` call at SubPlan.java:219 when the pipeline is already `CANCELING` looks cleaner to me than adding a separate `cancelRequested` flag — it also happens to resolve the FINISHED-vs-CANCELING interaction I raised as Issue 1 in my own review, since a pipeline that's still `CANCELING` would simply never get pushed to `FAILING` in the first place. @zhangshenghang, worth folding this into the same change: the existing E2E test is a good regression anchor for the "already mid-cancel" ordering, but per SEZ9's point it should also gain a case for a sibling vertex whose `cancel()` hasn't fired yet when the worker is reported lost, since that's the ordering the current fix doesn't close. I'm revising my own conclusion in light of this: the "recommended fix" I listed under Issue 1 is superseded by this more general problem. I'll re-review once a fix lands. -- 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]
