Skip to content

Commit de2697c

Browse files
[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
1 parent 0a59e04 commit de2697c

7 files changed

Lines changed: 304 additions & 8 deletions

File tree

streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/mapper/FlinkApplicationMapper.java

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,12 @@ public interface FlinkApplicationMapper extends BaseMapper<FlinkApplication> {
3333

3434
FlinkApplication selectApp(@Param("id") Long id);
3535

36-
void persistMetrics(@Param("app") FlinkApplication application);
36+
void persistMetrics(
37+
@Param("app") FlinkApplication application,
38+
@Param("cancellingState") int cancellingState,
39+
@Param("canceledState") int canceledState,
40+
@Param("cancellingOptionState") int cancellingOptionState,
41+
@Param("savepointingOptionState") int savepointingOptionState);
3742

3843
List<FlinkApplication> selectAppsByTeamId(@Param("teamId") Long teamId);
3944

streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/service/application/impl/FlinkApplicationManageServiceImpl.java

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -171,7 +171,12 @@ public void toEffective(FlinkApplication appParam) {
171171

172172
@Override
173173
public void persistMetrics(FlinkApplication appParam) {
174-
this.baseMapper.persistMetrics(appParam);
174+
this.baseMapper.persistMetrics(
175+
appParam,
176+
FlinkAppStateEnum.CANCELLING.getValue(),
177+
FlinkAppStateEnum.CANCELED.getValue(),
178+
OptionStateEnum.CANCELLING.getValue(),
179+
OptionStateEnum.SAVEPOINTING.getValue());
175180
}
176181

177182
@Override

streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/watcher/FlinkK8sChangeEventListener.java

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -188,8 +188,10 @@ private void setByJobStatusCV(FlinkApplication app, JobStatusCV jobStatus) {
188188
app.setStartTime(new Date(startTime > 0 ? startTime : 0));
189189
app.setEndTime(endTime > 0 && endTime >= startTime ? new Date(endTime) : null);
190190
app.setDuration(duration > 0 ? duration : 0);
191-
// when a flink job status change event can be received, it means
192-
// that the operation command sent by streampark has been completed.
193-
app.setOptionState(OptionStateEnum.NONE.getValue());
191+
// A non-terminal event can race with an in-flight cancellation. Keep the operation state
192+
// until the watcher observes a terminal job state.
193+
if (state != FlinkJobState.CANCELLING()) {
194+
app.setOptionState(OptionStateEnum.NONE.getValue());
195+
}
194196
}
195197
}

streampark-console/streampark-console-service/src/main/resources/mapper/core/FlinkApplicationMapper.xml

Lines changed: 35 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -124,7 +124,22 @@
124124
tracking=#{app.tracking},
125125
</if>
126126
<if test="app.optionState != null">
127-
option_state=#{app.optionState},
127+
<choose>
128+
<when test="@org.apache.streampark.console.core.enums.FlinkAppStateEnum@isEndState(app.state)">
129+
option_state=#{app.optionState},
130+
</when>
131+
<otherwise>
132+
option_state=case
133+
when state=#{cancellingState}
134+
and option_state in (
135+
#{cancellingOptionState},
136+
#{savepointingOptionState}
137+
)
138+
then option_state
139+
else #{app.optionState}
140+
end,
141+
</otherwise>
142+
</choose>
128143
</if>
129144
<if test="app.startTime != null">
130145
start_time=#{app.startTime},
@@ -165,9 +180,27 @@
165180
</if>
166181
</otherwise>
167182
</choose>
168-
state=#{app.state}
183+
<choose>
184+
<when test="@org.apache.streampark.console.core.enums.FlinkAppStateEnum@isEndState(app.state)">
185+
state=#{app.state},
186+
</when>
187+
<otherwise>
188+
state=case
189+
when state=#{cancellingState}
190+
and option_state in (
191+
#{cancellingOptionState},
192+
#{savepointingOptionState}
193+
)
194+
then state
195+
else #{app.state}
196+
end,
197+
</otherwise>
198+
</choose>
169199
</set>
170200
where id=#{app.id}
201+
<if test="@org.apache.streampark.console.core.enums.FlinkAppStateEnum@isEndState(app.state) == false">
202+
and state != #{canceledState}
203+
</if>
171204
</update>
172205

173206
<select id="selectAppsByTeamId" resultType="org.apache.streampark.console.core.entity.FlinkApplication" parameterType="java.lang.Long">

streampark-console/streampark-console-service/src/test/java/org/apache/streampark/console/core/service/FlinkApplicationManageServiceTest.java

Lines changed: 196 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,8 @@
2121
import org.apache.streampark.console.SpringUnitTestBase;
2222
import org.apache.streampark.console.core.entity.FlinkApplication;
2323
import org.apache.streampark.console.core.entity.YarnQueue;
24+
import org.apache.streampark.console.core.enums.FlinkAppStateEnum;
25+
import org.apache.streampark.console.core.enums.OptionStateEnum;
2426
import org.apache.streampark.console.core.service.application.FlinkApplicationActionService;
2527
import org.apache.streampark.console.core.service.application.FlinkApplicationManageService;
2628
import org.apache.streampark.console.core.service.application.impl.FlinkApplicationManageServiceImpl;
@@ -146,4 +148,198 @@ void testCheckQueueValidationIfNeeded() {
146148
app2.setYarnQueue(nonExistedQueue);
147149
assertThat(applicationServiceImpl.validateQueueIfNeeded(app1, app2)).isFalse();
148150
}
151+
152+
@Test
153+
void testPersistMetricsPreservesInFlightCancellation() {
154+
assertNonTerminalMetricsPreserveOperation(OptionStateEnum.CANCELLING);
155+
assertNonTerminalMetricsPreserveOperation(OptionStateEnum.SAVEPOINTING);
156+
}
157+
158+
@Test
159+
void testPersistMetricsRecoversInterruptedCancellation() {
160+
assertInterruptedCancellationRecovers(FlinkAppStateEnum.RUNNING);
161+
assertInterruptedCancellationRecovers(FlinkAppStateEnum.FAILING);
162+
}
163+
164+
@Test
165+
void testPersistMetricsUpdatesUnprotectedNonTerminalState() {
166+
FlinkApplication persisted = createApplication(
167+
FlinkAppStateEnum.RUNNING, OptionStateEnum.NONE);
168+
169+
FlinkApplication snapshot = new FlinkApplication();
170+
snapshot.setId(persisted.getId());
171+
snapshot.setState(FlinkAppStateEnum.FAILING.getValue());
172+
snapshot.setOptionState(OptionStateEnum.NONE.getValue());
173+
applicationManageService.persistMetrics(snapshot);
174+
175+
FlinkApplication actual = applicationManageService.getById(persisted.getId());
176+
assertThat(actual.getState()).isEqualTo(FlinkAppStateEnum.FAILING.getValue());
177+
assertThat(actual.getOptionState()).isEqualTo(OptionStateEnum.NONE.getValue());
178+
}
179+
180+
@Test
181+
void testPersistMetricsDoesNotBlockStandaloneSavepointCompletion() {
182+
FlinkApplication persisted = createApplication(
183+
FlinkAppStateEnum.RUNNING, OptionStateEnum.SAVEPOINTING);
184+
185+
FlinkApplication completedSnapshot = new FlinkApplication();
186+
completedSnapshot.setId(persisted.getId());
187+
completedSnapshot.setState(FlinkAppStateEnum.RUNNING.getValue());
188+
completedSnapshot.setOptionState(OptionStateEnum.NONE.getValue());
189+
applicationManageService.persistMetrics(completedSnapshot);
190+
191+
FlinkApplication actual = applicationManageService.getById(persisted.getId());
192+
assertThat(actual.getState()).isEqualTo(FlinkAppStateEnum.RUNNING.getValue());
193+
assertThat(actual.getOptionState()).isEqualTo(OptionStateEnum.NONE.getValue());
194+
}
195+
196+
@Test
197+
void testPersistMetricsConvergesCancellationOnTerminalState() {
198+
FlinkApplication persisted = createApplication(
199+
FlinkAppStateEnum.CANCELLING, OptionStateEnum.CANCELLING);
200+
201+
FlinkApplication terminalSnapshot = new FlinkApplication();
202+
terminalSnapshot.setId(persisted.getId());
203+
terminalSnapshot.setState(FlinkAppStateEnum.CANCELED.getValue());
204+
terminalSnapshot.setOptionState(OptionStateEnum.NONE.getValue());
205+
applicationManageService.persistMetrics(terminalSnapshot);
206+
207+
FlinkApplication actual = applicationManageService.getById(persisted.getId());
208+
assertThat(actual.getState()).isEqualTo(FlinkAppStateEnum.CANCELED.getValue());
209+
assertThat(actual.getOptionState()).isEqualTo(OptionStateEnum.NONE.getValue());
210+
}
211+
212+
@Test
213+
void testPersistMetricsRejectsLateNonTerminalSnapshotAfterCancellationCompletes() {
214+
assertCompletedCancellationRejectsLateSnapshot(
215+
FlinkAppStateEnum.FAILING, OptionStateEnum.NONE);
216+
assertCompletedCancellationRejectsLateSnapshot(
217+
FlinkAppStateEnum.RUNNING, OptionStateEnum.NONE);
218+
assertCompletedCancellationRejectsLateSnapshot(
219+
FlinkAppStateEnum.CANCELLING, OptionStateEnum.CANCELLING);
220+
}
221+
222+
@Test
223+
void testPersistMetricsAllowsExplicitRestartAfterCancellationCompletes() {
224+
FlinkApplication persisted = createApplication(
225+
FlinkAppStateEnum.CANCELED, OptionStateEnum.NONE);
226+
227+
FlinkApplication starting = new FlinkApplication();
228+
starting.setId(persisted.getId());
229+
starting.setState(FlinkAppStateEnum.STARTING.getValue());
230+
assertThat(applicationManageService.updateById(starting)).isTrue();
231+
232+
FlinkApplication runningSnapshot = new FlinkApplication();
233+
runningSnapshot.setId(persisted.getId());
234+
runningSnapshot.setState(FlinkAppStateEnum.RUNNING.getValue());
235+
runningSnapshot.setOptionState(OptionStateEnum.NONE.getValue());
236+
runningSnapshot.setTotalTask(2);
237+
applicationManageService.persistMetrics(runningSnapshot);
238+
239+
FlinkApplication actual = applicationManageService.getById(persisted.getId());
240+
assertThat(actual.getState()).isEqualTo(FlinkAppStateEnum.RUNNING.getValue());
241+
assertThat(actual.getOptionState()).isEqualTo(OptionStateEnum.NONE.getValue());
242+
assertThat(actual.getTotalTask()).isEqualTo(2);
243+
}
244+
245+
@Test
246+
void testPersistMetricsAllowsTerminalCorrectionAfterCancellationCompletes() {
247+
FlinkApplication persisted = createApplication(
248+
FlinkAppStateEnum.CANCELED, OptionStateEnum.NONE);
249+
250+
FlinkApplication failedSnapshot = new FlinkApplication();
251+
failedSnapshot.setId(persisted.getId());
252+
failedSnapshot.setState(FlinkAppStateEnum.FAILED.getValue());
253+
failedSnapshot.setOptionState(OptionStateEnum.NONE.getValue());
254+
applicationManageService.persistMetrics(failedSnapshot);
255+
256+
FlinkApplication actual = applicationManageService.getById(persisted.getId());
257+
assertThat(actual.getState()).isEqualTo(FlinkAppStateEnum.FAILED.getValue());
258+
assertThat(actual.getOptionState()).isEqualTo(OptionStateEnum.NONE.getValue());
259+
}
260+
261+
@Test
262+
void testPersistMetricsDoesNotFreezeRecoverableLostState() {
263+
FlinkApplication persisted = createApplication(
264+
FlinkAppStateEnum.LOST, OptionStateEnum.NONE);
265+
266+
FlinkApplication recoveredSnapshot = new FlinkApplication();
267+
recoveredSnapshot.setId(persisted.getId());
268+
recoveredSnapshot.setState(FlinkAppStateEnum.RUNNING.getValue());
269+
recoveredSnapshot.setOptionState(OptionStateEnum.NONE.getValue());
270+
applicationManageService.persistMetrics(recoveredSnapshot);
271+
272+
FlinkApplication actual = applicationManageService.getById(persisted.getId());
273+
assertThat(actual.getState()).isEqualTo(FlinkAppStateEnum.RUNNING.getValue());
274+
assertThat(actual.getOptionState()).isEqualTo(OptionStateEnum.NONE.getValue());
275+
}
276+
277+
private void assertCompletedCancellationRejectsLateSnapshot(
278+
FlinkAppStateEnum lateState,
279+
OptionStateEnum lateOptionState) {
280+
FlinkApplication persisted = createApplication(
281+
FlinkAppStateEnum.CANCELLING, OptionStateEnum.CANCELLING);
282+
283+
FlinkApplication terminalSnapshot = new FlinkApplication();
284+
terminalSnapshot.setId(persisted.getId());
285+
terminalSnapshot.setState(FlinkAppStateEnum.CANCELED.getValue());
286+
terminalSnapshot.setOptionState(OptionStateEnum.NONE.getValue());
287+
applicationManageService.persistMetrics(terminalSnapshot);
288+
289+
FlinkApplication lateSnapshot = new FlinkApplication();
290+
lateSnapshot.setId(persisted.getId());
291+
lateSnapshot.setState(lateState.getValue());
292+
lateSnapshot.setOptionState(lateOptionState.getValue());
293+
lateSnapshot.setTotalTask(7);
294+
applicationManageService.persistMetrics(lateSnapshot);
295+
296+
FlinkApplication actual = applicationManageService.getById(persisted.getId());
297+
assertThat(actual.getState()).isEqualTo(FlinkAppStateEnum.CANCELED.getValue());
298+
assertThat(actual.getOptionState()).isEqualTo(OptionStateEnum.NONE.getValue());
299+
assertThat(actual.getTotalTask()).isNull();
300+
}
301+
302+
private void assertInterruptedCancellationRecovers(FlinkAppStateEnum currentState) {
303+
FlinkApplication persisted = createApplication(
304+
FlinkAppStateEnum.CANCELLING, OptionStateEnum.NONE);
305+
306+
FlinkApplication currentSnapshot = new FlinkApplication();
307+
currentSnapshot.setId(persisted.getId());
308+
currentSnapshot.setState(currentState.getValue());
309+
currentSnapshot.setOptionState(OptionStateEnum.NONE.getValue());
310+
currentSnapshot.setTotalTask(4);
311+
applicationManageService.persistMetrics(currentSnapshot);
312+
313+
FlinkApplication actual = applicationManageService.getById(persisted.getId());
314+
assertThat(actual.getState()).isEqualTo(currentState.getValue());
315+
assertThat(actual.getOptionState()).isEqualTo(OptionStateEnum.NONE.getValue());
316+
assertThat(actual.getTotalTask()).isEqualTo(4);
317+
}
318+
319+
private void assertNonTerminalMetricsPreserveOperation(OptionStateEnum optionState) {
320+
FlinkApplication persisted = createApplication(FlinkAppStateEnum.CANCELLING, optionState);
321+
322+
FlinkApplication staleSnapshot = new FlinkApplication();
323+
staleSnapshot.setId(persisted.getId());
324+
staleSnapshot.setState(FlinkAppStateEnum.FAILING.getValue());
325+
staleSnapshot.setOptionState(OptionStateEnum.NONE.getValue());
326+
staleSnapshot.setTotalTask(3);
327+
applicationManageService.persistMetrics(staleSnapshot);
328+
329+
FlinkApplication actual = applicationManageService.getById(persisted.getId());
330+
assertThat(actual.getState()).isEqualTo(FlinkAppStateEnum.CANCELLING.getValue());
331+
assertThat(actual.getOptionState()).isEqualTo(optionState.getValue());
332+
assertThat(actual.getTotalTask()).isEqualTo(3);
333+
}
334+
335+
private FlinkApplication createApplication(
336+
FlinkAppStateEnum state,
337+
OptionStateEnum optionState) {
338+
FlinkApplication application = new FlinkApplication();
339+
application.setTeamId(1L);
340+
application.setState(state.getValue());
341+
application.setOptionState(optionState.getValue());
342+
assertThat(applicationManageService.save(application)).isTrue();
343+
return application;
344+
}
149345
}

streampark-console/streampark-console-webapp/src/views/flink/app/hooks/useAppTableAction.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -104,7 +104,8 @@ export const useAppTableAction = (
104104
class: 'e2e-flinkapp-cancel-btn',
105105
tooltip: { title: t('flink.app.operation.cancel') },
106106
ifShow:
107-
record.state == AppStateEnum.RUNNING && record['optionState'] == OptionStateEnum.NONE,
107+
[AppStateEnum.RUNNING, AppStateEnum.FAILING].includes(record.state) &&
108+
record['optionState'] == OptionStateEnum.NONE,
108109
auth: 'app:cancel',
109110
icon: 'ant-design:pause-circle-outlined',
110111
onClick: handleCancel.bind(null, record),
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one or more
3+
* contributor license agreements. See the NOTICE file distributed with
4+
* this work for additional information regarding copyright ownership.
5+
* The ASF licenses this file to You under the Apache License, Version 2.0
6+
* (the "License"); you may not use this file except in compliance with
7+
* the License. You may obtain a copy of the License at
8+
*
9+
* http://www.apache.org/licenses/LICENSE-2.0
10+
*
11+
* Unless required by applicable law or agreed to in writing, software
12+
* distributed under the License is distributed on an "AS IS" BASIS,
13+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14+
* See the License for the specific language governing permissions and
15+
* limitations under the License.
16+
*/
17+
18+
package org.apache.streampark.flink.kubernetes
19+
20+
import org.apache.streampark.flink.kubernetes.enums.FlinkJobState
21+
import org.apache.streampark.flink.kubernetes.watcher.FlinkJobStatusWatcher
22+
23+
import org.junit.jupiter.api.Assertions.assertEquals
24+
import org.junit.jupiter.api.Test
25+
26+
class FlinkJobStatusWatcherTest {
27+
28+
@Test
29+
def doNotInferActiveCancellationFromPersistedState(): Unit = {
30+
Seq(FlinkJobState.RUNNING, FlinkJobState.FAILING, FlinkJobState.SILENT).foreach {
31+
current =>
32+
assertEquals(
33+
current,
34+
FlinkJobStatusWatcher.inferFlinkJobStateFromPersist(
35+
current,
36+
FlinkJobState.CANCELLING))
37+
}
38+
}
39+
40+
@Test
41+
def convergeCancellationWhenWatcherObservesTerminalState(): Unit = {
42+
assertEquals(
43+
FlinkJobState.CANCELED,
44+
FlinkJobStatusWatcher.inferFlinkJobStateFromPersist(
45+
FlinkJobState.TERMINATED,
46+
FlinkJobState.CANCELLING))
47+
assertEquals(
48+
FlinkJobState.FAILED,
49+
FlinkJobStatusWatcher.inferFlinkJobStateFromPersist(
50+
FlinkJobState.FAILED,
51+
FlinkJobState.CANCELLING))
52+
}
53+
54+
}

0 commit comments

Comments
 (0)