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.

Reply via email to