DanielLeens commented on PR #12311: URL: https://github.com/apache/seatunnel/pull/12311#issuecomment-5998161764
@SEZ9 The revision I owed you is up: `f85d43600df`. It covers the points from your review and your follow-up. Point by point: 1. **Engine-initiated cancel.** `resolveLostMemberState` now takes the job status as well as the vertex state, and a vertex only resolves to `CANCELED` while the job itself is `JobStatus.CANCELING` (set by `PhysicalPlan#cancelJob` / `#stopJob`). In every other job status a `CANCELING` vertex resolves to `FAILED` again, as on `dev`. That covers `SubPlan#restorePipelineState` after a master switch, `SubPlan#handleCheckpointError`, and a `FAILING` pipeline or job cancelling its siblings, so a job the user never cancelled cannot end `CANCELED` with the reason dropped. 2. **Sibling vertices still `RUNNING`.** While the job is `CANCELING`, `DEPLOYING` and `RUNNING` vertices on the lost member now resolve to `CANCELED` too, so the vertices the sequential cancel loop has not reached yet no longer flip the job to `FAILED`. I used the job status rather than the pipeline status because the pipeline status cannot tell a user cancel from an engine cancel (point 1). The state read, the decision and the transition now run inside `synchronized (physicalVertex)`, the same monitor `updateTaskState`, `cancel` and `stateProcess` use, and the job status is read per vertex inside that block. 3. **IT Javadoc.** Reworded to state what is now true: every vertex of the user-cancelled job that is still deployed on the lost worker reaches `CANCELED`, and an engine-cancelled vertex keeps resolving to `FAILED`. 4. **Dropped reason.** `resolveTasksOnLostMember` logs a WARN with the task group, the lost address, the job status and the previous state whenever a member loss completes a cancel. 5. **Restore argument.** The Javadoc now names the real guard: a pipeline that turned `FAILING` without a failed task stays `FAILED` because `SubPlan#getPipelineEndState` checks the `FAILING` pipeline status, not because of `failedTaskNum`. 6. **Naming.** The private `makeTasksFailed` is now `resolveTasksOnLostMember`, and `resolveLostMemberState` returns `Optional<ExecutionState>` instead of a `null` sentinel. I kept the public `failedTaskOnMemberRemoved` name and gave it a Javadoc that says it can also resolve to `CANCELED`; the other IT Javadocs that cite it describe the non-cancelled `FAILED` path, which is unchanged, so I left them alone. `CoordinatorServiceLostMemberResolutionTest` now pins the whole matrix: `CANCELING`/`DEPLOYING`/`RUNNING` under a `CANCELING` job resolve to `CANCELED`, the same three states under every other job status (and a missing one) resolve to `FAILED`, and all other vertex states are left alone. One limit to be clear about: the unit test covers the decision, not the locking. The end-to-end behaviour is exercised by `SplitClusterFaultToleranceIT#testStreamJobCancelResolvesWhenWorkerCrashesBeforeCancelAck` in `engine-v2-it`. I have not compiled or run anything locally beyond `spotless:apply`, so CI on this head is the verification; I will post the result once it has run. @davidzollo your approval was on the previous head, so this needs a fresh look as well. -- 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]
