This is an automated email from the ASF dual-hosted git repository.
gyfora pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/flink-kubernetes-operator.git
The following commit(s) were added to refs/heads/main by this push:
new e3db28e6 [FLINK-40374] Fix redeploy-from-savepoint for suspended jobs
(#1176)
e3db28e6 is described below
commit e3db28e6522656c92d8af0da6cf6577ab529da7b
Author: Dale Lane <[email protected]>
AuthorDate: Thu Aug 13 09:02:34 2026 +0100
[FLINK-40374] Fix redeploy-from-savepoint for suspended jobs (#1176)
---
.../operator/service/AbstractFlinkService.java | 20 ++++++++++--
.../sessionjob/SessionJobReconcilerTest.java | 36 ++++++++++++++++++++++
.../operator/service/AbstractFlinkServiceTest.java | 28 +++++++++++++++++
3 files changed, 81 insertions(+), 3 deletions(-)
diff --git
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java
index bc0e9ccf..ad650423 100644
---
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java
+++
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java
@@ -379,9 +379,13 @@ public abstract class AbstractFlinkService implements
FlinkService {
try (var clusterClient = getClusterClient(conf)) {
switch (suspendMode) {
case STATELESS:
+ if (!cancelJobOrError(clusterClient, status, true)) {
+ // This is async we need to return and re-observe
+ return CancelResult.pending();
+ }
+ break;
case CANCEL:
- if (!cancelJobOrError(
- clusterClient, status, suspendMode ==
SuspendMode.STATELESS)) {
+ if (!cancelJobOrError(clusterClient, status, false)) {
// This is async we need to return and re-observe
return CancelResult.pending();
}
@@ -416,7 +420,17 @@ public abstract class AbstractFlinkService implements
FlinkService {
RestClusterClient<String> clusterClient,
CommonStatus<?> status,
boolean ignoreMissing) {
- var jobID = JobID.fromHexString(status.getJobStatus().getJobId());
+ var jobIdString = status.getJobStatus().getJobId();
+ if (jobIdString == null) {
+ if (ignoreMissing) {
+ LOG.info("Job already missing");
+ return true;
+ }
+ throw new UpgradeFailureException(
+ "Cannot find job when trying to cancel",
+ EventRecorder.Reason.CleanupFailed.name());
+ }
+ var jobID = JobID.fromHexString(jobIdString);
if (ReconciliationUtils.isJobCancelling(status)) {
LOG.info("Job already cancelling");
return false;
diff --git
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconcilerTest.java
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconcilerTest.java
index 3416d474..12e4550c 100644
---
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconcilerTest.java
+++
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconcilerTest.java
@@ -1091,4 +1091,40 @@ public class SessionJobReconcilerTest extends
OperatorTestBase {
assertFalse(deleteControl.isRemoveFinalizer());
assertEquals(10_000, deleteControl.getScheduleDelay().get());
}
+
+ @Test
+ public void testSavepointRedeployAfterSavepointSuspend() throws Exception {
+ // submit job
+ var readyCtx =
TestUtils.createContextWithReadyFlinkDeployment(kubernetesClient);
+ FlinkSessionJob sessionJob = TestUtils.buildSessionJob();
+ sessionJob.getSpec().getJob().setUpgradeMode(UpgradeMode.SAVEPOINT);
+
+ // reconcile and verify
+ reconciler.reconcile(sessionJob, readyCtx);
+ verifyAndSetRunningJobsToStatus(
+ sessionJob, JobState.RUNNING, RECONCILING, null,
flinkService.listJobs());
+
Assertions.assertNotNull(sessionJob.getStatus().getJobStatus().getJobId());
+
+ // savepoint suspend
+ sessionJob.getSpec().getJob().setState(JobState.SUSPENDED);
+
+ // reconcile and verify
+ reconciler.reconcile(sessionJob, readyCtx);
+ verifyJobState(sessionJob, JobState.SUSPENDED, FINISHED);
+ assertNull(sessionJob.getStatus().getJobStatus().getJobId());
+
+ flinkService.clear();
+
+ // request redeploy from explicit savepoint
+ sessionJob.getSpec().getJob().setState(JobState.RUNNING);
+
sessionJob.getSpec().getJob().setInitialSavepointPath("s3://bucket/savepoint-explicit");
+ sessionJob.getSpec().getJob().setSavepointRedeployNonce(1L);
+
+ // reconcile and verify
+ reconciler.reconcile(sessionJob, readyCtx);
+ assertEquals(1, flinkService.listJobs().size());
+ assertEquals(
+ "s3://bucket/savepoint-explicit",
+ verifyAndReturnTheSubmittedJob(sessionJob,
flinkService.listJobs()).f0);
+ }
}
diff --git
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkServiceTest.java
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkServiceTest.java
index a6f3ed69..bb00337c 100644
---
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkServiceTest.java
+++
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkServiceTest.java
@@ -388,6 +388,34 @@ public class AbstractFlinkServiceTest {
assertNull(jobStatus.getJobId());
}
+ /**
+ * A session job suspended with a savepoint has its jobId cleared and a
terminal state.
+ * Redeploying from that state should recognise the null jobId as an
already-missing job instead
+ * of trying to parse it.
+ */
+ @Test
+ public void cancelSessionJobWithStatelessModeAndNoJobId() throws Exception
{
+ var testingClusterClient =
+ new TestingClusterClient<>(configuration,
TestUtils.TEST_DEPLOYMENT_NAME);
+ testingClusterClient.setCancelFunction(
+ jobID -> {
+ fail("cancel should not be called for a job with no
recorded jobId");
+ return null;
+ });
+ var flinkService = new TestingService(testingClusterClient);
+
+ var job = TestUtils.buildSessionJob();
+ var jobStatus = job.getStatus().getJobStatus();
+ jobStatus.setJobId(null);
+ jobStatus.setState(FINISHED);
+ ReconciliationUtils.updateStatusForDeployedSpec(job, new
Configuration());
+
+ var result = flinkService.cancelSessionJob(job, SuspendMode.STATELESS,
new Configuration());
+ assertFalse(result.isPending());
+ assertEquals(FINISHED, jobStatus.getState());
+ assertNull(jobStatus.getJobId());
+ }
+
/**
* Reproduces the operator-upgrade scenario for Session Mode with CANCEL
upgrade mode: when a
* running session job's JobManager has already moved the job into a
terminal state (e.g.