From de2697c9e5bc350e7bdb466818ff6457f4636b1f Mon Sep 17 00:00:00 2001 From: Power John Date: Sat, 25 Jul 2026 23:57:07 +0800 Subject: [PATCH] [K8s] Fix failed application cancellation state race Allow failed Kubernetes applications to be canceled while protecting active and completed cancellation from stale watcher snapshots. Keep the database guard scoped so interrupted cancellation recovers after Console restart, explicit restarts and terminal corrections remain authoritative, and standalone savepoint completion remains unaffected. Refs #4332 --- .../core/mapper/FlinkApplicationMapper.java | 7 +- .../FlinkApplicationManageServiceImpl.java | 7 +- .../watcher/FlinkK8sChangeEventListener.java | 8 +- .../mapper/core/FlinkApplicationMapper.xml | 37 +++- .../FlinkApplicationManageServiceTest.java | 196 ++++++++++++++++++ .../flink/app/hooks/useAppTableAction.ts | 3 +- .../FlinkJobStatusWatcherTest.scala | 54 +++++ 7 files changed, 304 insertions(+), 8 deletions(-) create mode 100644 streampark-flink/streampark-flink-kubernetes/src/test/scala/org/apache/streampark/flink/kubernetes/FlinkJobStatusWatcherTest.scala diff --git a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/mapper/FlinkApplicationMapper.java b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/mapper/FlinkApplicationMapper.java index 3338b090ca..5b5afc3bb5 100644 --- a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/mapper/FlinkApplicationMapper.java +++ b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/mapper/FlinkApplicationMapper.java @@ -33,7 +33,12 @@ public interface FlinkApplicationMapper extends BaseMapper { FlinkApplication selectApp(@Param("id") Long id); - void persistMetrics(@Param("app") FlinkApplication application); + void persistMetrics( + @Param("app") FlinkApplication application, + @Param("cancellingState") int cancellingState, + @Param("canceledState") int canceledState, + @Param("cancellingOptionState") int cancellingOptionState, + @Param("savepointingOptionState") int savepointingOptionState); List selectAppsByTeamId(@Param("teamId") Long teamId); diff --git a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/service/application/impl/FlinkApplicationManageServiceImpl.java b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/service/application/impl/FlinkApplicationManageServiceImpl.java index 0d18aab3d2..545d3514f7 100644 --- a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/service/application/impl/FlinkApplicationManageServiceImpl.java +++ b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/service/application/impl/FlinkApplicationManageServiceImpl.java @@ -171,7 +171,12 @@ public void toEffective(FlinkApplication appParam) { @Override public void persistMetrics(FlinkApplication appParam) { - this.baseMapper.persistMetrics(appParam); + this.baseMapper.persistMetrics( + appParam, + FlinkAppStateEnum.CANCELLING.getValue(), + FlinkAppStateEnum.CANCELED.getValue(), + OptionStateEnum.CANCELLING.getValue(), + OptionStateEnum.SAVEPOINTING.getValue()); } @Override diff --git a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/watcher/FlinkK8sChangeEventListener.java b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/watcher/FlinkK8sChangeEventListener.java index 8a57a07a17..f9379aefa8 100644 --- a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/watcher/FlinkK8sChangeEventListener.java +++ b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/watcher/FlinkK8sChangeEventListener.java @@ -188,8 +188,10 @@ private void setByJobStatusCV(FlinkApplication app, JobStatusCV jobStatus) { app.setStartTime(new Date(startTime > 0 ? startTime : 0)); app.setEndTime(endTime > 0 && endTime >= startTime ? new Date(endTime) : null); app.setDuration(duration > 0 ? duration : 0); - // when a flink job status change event can be received, it means - // that the operation command sent by streampark has been completed. - app.setOptionState(OptionStateEnum.NONE.getValue()); + // A non-terminal event can race with an in-flight cancellation. Keep the operation state + // until the watcher observes a terminal job state. + if (state != FlinkJobState.CANCELLING()) { + app.setOptionState(OptionStateEnum.NONE.getValue()); + } } } diff --git a/streampark-console/streampark-console-service/src/main/resources/mapper/core/FlinkApplicationMapper.xml b/streampark-console/streampark-console-service/src/main/resources/mapper/core/FlinkApplicationMapper.xml index 0c567e6d0d..103d5a0195 100644 --- a/streampark-console/streampark-console-service/src/main/resources/mapper/core/FlinkApplicationMapper.xml +++ b/streampark-console/streampark-console-service/src/main/resources/mapper/core/FlinkApplicationMapper.xml @@ -124,7 +124,22 @@ tracking=#{app.tracking}, - option_state=#{app.optionState}, + + + option_state=#{app.optionState}, + + + option_state=case + when state=#{cancellingState} + and option_state in ( + #{cancellingOptionState}, + #{savepointingOptionState} + ) + then option_state + else #{app.optionState} + end, + + start_time=#{app.startTime}, @@ -165,9 +180,27 @@ - state=#{app.state} + + + state=#{app.state}, + + + state=case + when state=#{cancellingState} + and option_state in ( + #{cancellingOptionState}, + #{savepointingOptionState} + ) + then state + else #{app.state} + end, + + where id=#{app.id} + + and state != #{canceledState} +