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} +