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 15f648ce [FLINK-35776] Simplify job status handling
15f648ce is described below

commit 15f648ce46537d6a1df3a87cfd653a3f855d0dcd
Author: Gyula Fora <[email protected]>
AuthorDate: Thu May 23 16:47:24 2024 +0200

    [FLINK-35776] Simplify job status handling
---
 .../flink/autoscaler/utils/JobStatusUtils.java     |  16 +--
 .../exception/MissingSessionJobException.java      |  33 -------
 .../operator/exception/UnknownJobException.java    |  33 -------
 .../operator/observer/JobStatusObserver.java       |  94 +++++++-----------
 .../observer/deployment/ApplicationObserver.java   |  49 ----------
 .../observer/deployment/SessionObserver.java       |   2 +-
 .../sessionjob/FlinkSessionJobObserver.java        | 107 +++------------------
 .../deployment/ApplicationReconciler.java          |   8 +-
 .../sessionjob/SessionJobReconciler.java           |  18 +---
 .../operator/service/AbstractFlinkService.java     |  61 +++++++++---
 .../kubernetes/operator/service/FlinkService.java  |   3 +-
 .../kubernetes/operator/TestingFlinkService.java   |  10 +-
 .../deployment/ApplicationObserverTest.java        |  25 ++---
 .../sessionjob/FlinkSessionJobObserverTest.java    |  24 +----
 .../operator/service/AbstractFlinkServiceTest.java |  36 +++++++
 15 files changed, 174 insertions(+), 345 deletions(-)

diff --git 
a/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/utils/JobStatusUtils.java
 
b/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/utils/JobStatusUtils.java
index c4b99231..9611088b 100644
--- 
a/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/utils/JobStatusUtils.java
+++ 
b/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/utils/JobStatusUtils.java
@@ -38,16 +38,18 @@ public class JobStatusUtils {
     public static List<JobStatusMessage> toJobStatusMessage(
             MultipleJobsDetails multipleJobsDetails) {
         return multipleJobsDetails.getJobs().stream()
-                .map(
-                        details ->
-                                new JobStatusMessage(
-                                        details.getJobId(),
-                                        details.getJobName(),
-                                        getEffectiveStatus(details),
-                                        details.getStartTime()))
+                .map(JobStatusUtils::toJobStatusMessage)
                 .collect(Collectors.toList());
     }
 
+    public static JobStatusMessage toJobStatusMessage(JobDetails details) {
+        return new JobStatusMessage(
+                details.getJobId(),
+                details.getJobName(),
+                getEffectiveStatus(details),
+                details.getStartTime());
+    }
+
     @VisibleForTesting
     static JobStatus getEffectiveStatus(JobDetails details) {
         int numRunning = 
details.getTasksPerState()[ExecutionState.RUNNING.ordinal()];
diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/exception/MissingSessionJobException.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/exception/MissingSessionJobException.java
deleted file mode 100644
index ab926195..00000000
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/exception/MissingSessionJobException.java
+++ /dev/null
@@ -1,33 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements.  See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.flink.kubernetes.operator.exception;
-
-/** Exception to signal missing session job. */
-public class MissingSessionJobException extends RuntimeException {
-    public MissingSessionJobException(Throwable cause) {
-        super(cause);
-    }
-
-    public MissingSessionJobException(String msg) {
-        super(msg);
-    }
-
-    public MissingSessionJobException(String msg, Throwable cause) {
-        super(msg, cause);
-    }
-}
diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/exception/UnknownJobException.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/exception/UnknownJobException.java
deleted file mode 100644
index 34bff05e..00000000
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/exception/UnknownJobException.java
+++ /dev/null
@@ -1,33 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements.  See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.flink.kubernetes.operator.exception;
-
-/** Exception to signal unrecognized job found. */
-public class UnknownJobException extends RuntimeException {
-    public UnknownJobException(Throwable cause) {
-        super(cause);
-    }
-
-    public UnknownJobException(String msg) {
-        super(msg);
-    }
-
-    public UnknownJobException(String msg, Throwable cause) {
-        super(msg, cause);
-    }
-}
diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/JobStatusObserver.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/JobStatusObserver.java
index d30cbfda..08f74f79 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/JobStatusObserver.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/JobStatusObserver.java
@@ -17,6 +17,7 @@
 
 package org.apache.flink.kubernetes.operator.observer;
 
+import org.apache.flink.api.common.JobID;
 import org.apache.flink.kubernetes.operator.api.AbstractFlinkResource;
 import org.apache.flink.kubernetes.operator.api.spec.JobState;
 import org.apache.flink.kubernetes.operator.api.status.JobStatus;
@@ -24,23 +25,21 @@ import 
org.apache.flink.kubernetes.operator.controller.FlinkResourceContext;
 import org.apache.flink.kubernetes.operator.reconciler.ReconciliationUtils;
 import org.apache.flink.kubernetes.operator.utils.EventRecorder;
 import org.apache.flink.runtime.client.JobStatusMessage;
+import org.apache.flink.runtime.rest.NotFoundException;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
-import java.util.ArrayList;
-import java.util.List;
-import java.util.Optional;
 import java.util.concurrent.TimeoutException;
 
 import static 
org.apache.flink.kubernetes.operator.utils.FlinkResourceExceptionUtils.updateFlinkResourceException;
 
 /** An observer to observe the job status. */
-public abstract class JobStatusObserver<R extends AbstractFlinkResource<?, ?>> 
{
+public class JobStatusObserver<R extends AbstractFlinkResource<?, ?>> {
 
     private static final Logger LOG = 
LoggerFactory.getLogger(JobStatusObserver.class);
 
-    public static final String MISSING_SESSION_JOB_ERR = "Missing Session Job";
+    public static final String JOB_NOT_FOUND_ERR = "Job Not Found";
 
     protected final EventRecorder eventRecorder;
 
@@ -68,43 +67,29 @@ public abstract class JobStatusObserver<R extends 
AbstractFlinkResource<?, ?>> {
         LOG.debug("Observing job status");
         var previousJobStatus = jobStatus.getState();
 
-        List<JobStatusMessage> clusterJobStatuses;
         try {
-            // Query job list from the cluster
-            clusterJobStatuses =
-                    new 
ArrayList<>(ctx.getFlinkService().listJobs(ctx.getObserveConfig()));
+            var newJobStatusOpt =
+                    ctx.getFlinkService()
+                            .getJobStatus(
+                                    ctx.getObserveConfig(),
+                                    JobID.fromHexString(jobStatus.getJobId()));
+
+            if (newJobStatusOpt.isPresent()) {
+                updateJobStatus(ctx, newJobStatusOpt.get());
+                
ReconciliationUtils.checkAndUpdateStableSpec(resource.getStatus());
+                return true;
+            } else {
+                onTargetJobNotFound(ctx);
+            }
         } catch (Exception e) {
             // Error while accessing the rest api, will try again later...
-            LOG.warn("Exception while listing jobs", e);
+            LOG.warn("Exception while getting job status", e);
             ifRunningMoveToReconciling(jobStatus, previousJobStatus);
             if (e instanceof TimeoutException) {
                 onTimeout(ctx);
             }
-            return false;
-        }
-
-        if (!clusterJobStatuses.isEmpty()) {
-            // There are jobs on the cluster, we filter the ones for this 
resource
-            Optional<JobStatusMessage> targetJobStatusMessage =
-                    filterTargetJob(jobStatus, clusterJobStatuses);
-
-            if (targetJobStatusMessage.isEmpty()) {
-                LOG.warn("No matching jobs found on the cluster");
-                ifRunningMoveToReconciling(jobStatus, previousJobStatus);
-                onTargetJobNotFound(ctx);
-                return false;
-            } else {
-                updateJobStatus(ctx, targetJobStatusMessage.get());
-            }
-            ReconciliationUtils.checkAndUpdateStableSpec(resource.getStatus());
-            return true;
-        } else {
-            LOG.debug("No jobs found on the cluster");
-            // No jobs found on the cluster, it is possible that the 
jobmanager is still starting up
-            ifRunningMoveToReconciling(jobStatus, previousJobStatus);
-            onNoJobsFound(ctx);
-            return false;
         }
+        return false;
     }
 
     /**
@@ -112,14 +97,21 @@ public abstract class JobStatusObserver<R extends 
AbstractFlinkResource<?, ?>> {
      *
      * @param ctx The Flink resource context.
      */
-    protected abstract void onTargetJobNotFound(FlinkResourceContext<R> ctx);
-
-    /**
-     * Callback when no jobs were found on the cluster.
-     *
-     * @param ctx The Flink resource context.
-     */
-    protected void onNoJobsFound(FlinkResourceContext<R> ctx) {}
+    protected void onTargetJobNotFound(FlinkResourceContext<R> ctx) {
+        ctx.getResource()
+                .getStatus()
+                .getJobStatus()
+                
.setState(org.apache.flink.api.common.JobStatus.RECONCILING.name());
+        ReconciliationUtils.updateForReconciliationError(
+                ctx, new NotFoundException(JOB_NOT_FOUND_ERR));
+        eventRecorder.triggerEvent(
+                ctx.getResource(),
+                EventRecorder.Type.Warning,
+                EventRecorder.Reason.Missing,
+                EventRecorder.Component.Job,
+                JOB_NOT_FOUND_ERR,
+                ctx.getKubernetesClient());
+    }
 
     /**
      * If we observed the job previously in RUNNING state we move to 
RECONCILING instead as we are
@@ -139,18 +131,7 @@ public abstract class JobStatusObserver<R extends 
AbstractFlinkResource<?, ?>> {
      *
      * @param ctx Resource context.
      */
-    protected abstract void onTimeout(FlinkResourceContext<R> ctx);
-
-    /**
-     * Filter the target job status message by the job list from the cluster.
-     *
-     * @param status the target job status.
-     * @param clusterJobStatuses the candidate cluster jobs.
-     * @return The target job status message. If no matched job found, {@code 
Optional.empty()} will
-     *     be returned.
-     */
-    protected abstract Optional<JobStatusMessage> filterTargetJob(
-            JobStatus status, List<JobStatusMessage> clusterJobStatuses);
+    protected void onTimeout(FlinkResourceContext<R> ctx) {}
 
     /**
      * Update the status in CR according to the cluster job status.
@@ -161,16 +142,13 @@ public abstract class JobStatusObserver<R extends 
AbstractFlinkResource<?, ?>> {
     private void updateJobStatus(FlinkResourceContext<R> ctx, JobStatusMessage 
clusterJobStatus) {
         var resource = ctx.getResource();
         var jobStatus = resource.getStatus().getJobStatus();
-        var previousJobId = jobStatus.getJobId();
         var previousJobStatus = jobStatus.getState();
 
         jobStatus.setState(clusterJobStatus.getJobState().name());
         jobStatus.setJobName(clusterJobStatus.getJobName());
-        jobStatus.setJobId(clusterJobStatus.getJobId().toHexString());
         
jobStatus.setStartTime(String.valueOf(clusterJobStatus.getStartTime()));
 
-        if (jobStatus.getJobId().equals(previousJobId)
-                && jobStatus.getState().equals(previousJobStatus)) {
+        if (jobStatus.getState().equals(previousJobStatus)) {
             LOG.debug("Job status ({}) unchanged", previousJobStatus);
         } else {
             
jobStatus.setUpdateTime(String.valueOf(System.currentTimeMillis()));
diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/deployment/ApplicationObserver.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/deployment/ApplicationObserver.java
index ba4002e1..90555407 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/deployment/ApplicationObserver.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/deployment/ApplicationObserver.java
@@ -19,19 +19,11 @@ package 
org.apache.flink.kubernetes.operator.observer.deployment;
 
 import org.apache.flink.kubernetes.operator.api.FlinkDeployment;
 import org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentStatus;
-import org.apache.flink.kubernetes.operator.api.status.JobStatus;
 import org.apache.flink.kubernetes.operator.controller.FlinkResourceContext;
-import org.apache.flink.kubernetes.operator.exception.UnknownJobException;
 import org.apache.flink.kubernetes.operator.observer.ClusterHealthObserver;
 import org.apache.flink.kubernetes.operator.observer.JobStatusObserver;
 import org.apache.flink.kubernetes.operator.observer.SnapshotObserver;
-import org.apache.flink.kubernetes.operator.reconciler.ReconciliationUtils;
 import org.apache.flink.kubernetes.operator.utils.EventRecorder;
-import org.apache.flink.runtime.client.JobStatusMessage;
-
-import java.util.Comparator;
-import java.util.List;
-import java.util.Optional;
 
 import static 
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions.OPERATOR_CLUSTER_HEALTH_CHECK_ENABLED;
 
@@ -73,46 +65,5 @@ public class ApplicationObserver extends 
AbstractFlinkDeploymentObserver {
         public void onTimeout(FlinkResourceContext<FlinkDeployment> ctx) {
             observeJmDeployment(ctx);
         }
-
-        @Override
-        protected Optional<JobStatusMessage> filterTargetJob(
-                JobStatus status, List<JobStatusMessage> clusterJobStatuses) {
-            if (!clusterJobStatuses.isEmpty()) {
-                clusterJobStatuses.sort(
-                        
Comparator.comparingLong(JobStatusMessage::getStartTime).reversed());
-                return Optional.of(clusterJobStatuses.get(0));
-            }
-            return Optional.empty();
-        }
-
-        @Override
-        protected void 
onTargetJobNotFound(FlinkResourceContext<FlinkDeployment> ctx) {
-            // This should never happen for application clusters, there is 
something
-            // wrong
-            setUnknownJobError(ctx);
-        }
-
-        /**
-         * We found a job on an application cluster that doesn't match the 
expected job. Trigger
-         * error.
-         *
-         * @param ctx Application deployment context.
-         */
-        private void setUnknownJobError(FlinkResourceContext<FlinkDeployment> 
ctx) {
-            ctx.getResource()
-                    .getStatus()
-                    .getJobStatus()
-                    
.setState(org.apache.flink.api.common.JobStatus.RECONCILING.name());
-            String err = "Unrecognized Job for Application deployment";
-            logger.error(err);
-            ReconciliationUtils.updateForReconciliationError(ctx, new 
UnknownJobException(err));
-            eventRecorder.triggerEvent(
-                    ctx.getResource(),
-                    EventRecorder.Type.Warning,
-                    EventRecorder.Reason.Missing,
-                    EventRecorder.Component.Job,
-                    err,
-                    ctx.getKubernetesClient());
-        }
     }
 }
diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/deployment/SessionObserver.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/deployment/SessionObserver.java
index f1621b1d..2d716bfd 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/deployment/SessionObserver.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/deployment/SessionObserver.java
@@ -36,7 +36,7 @@ public class SessionObserver extends 
AbstractFlinkDeploymentObserver {
         // Check if session cluster can serve rest calls following our 
practice in JobObserver
         try {
             logger.debug("Observing session cluster");
-            ctx.getFlinkService().listJobs(ctx.getObserveConfig());
+            ctx.getFlinkService().getClusterInfo(ctx.getObserveConfig());
             var rs = ctx.getResource().getStatus().getReconciliationStatus();
             if (rs.getState() == ReconciliationState.DEPLOYED) {
                 rs.markReconciledSpecAsStable();
diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/sessionjob/FlinkSessionJobObserver.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/sessionjob/FlinkSessionJobObserver.java
index 1846b156..d02e677b 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/sessionjob/FlinkSessionJobObserver.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/sessionjob/FlinkSessionJobObserver.java
@@ -17,39 +17,29 @@
 
 package org.apache.flink.kubernetes.operator.observer.sessionjob;
 
+import org.apache.flink.api.common.JobID;
 import org.apache.flink.kubernetes.operator.api.FlinkSessionJob;
 import org.apache.flink.kubernetes.operator.api.status.FlinkSessionJobStatus;
-import org.apache.flink.kubernetes.operator.api.status.JobStatus;
 import org.apache.flink.kubernetes.operator.controller.FlinkResourceContext;
-import 
org.apache.flink.kubernetes.operator.exception.MissingSessionJobException;
 import 
org.apache.flink.kubernetes.operator.observer.AbstractFlinkResourceObserver;
 import org.apache.flink.kubernetes.operator.observer.JobStatusObserver;
 import org.apache.flink.kubernetes.operator.observer.SnapshotObserver;
-import org.apache.flink.kubernetes.operator.reconciler.ReconciliationUtils;
 import org.apache.flink.kubernetes.operator.utils.EventRecorder;
-import org.apache.flink.runtime.client.JobStatusMessage;
-import org.apache.flink.runtime.jobmanager.HighAvailabilityMode;
-import org.apache.flink.util.Preconditions;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
-import java.util.Collection;
-import java.util.List;
-import java.util.Optional;
-import java.util.stream.Collectors;
-
 /** The observer of {@link FlinkSessionJob}. */
 public class FlinkSessionJobObserver extends 
AbstractFlinkResourceObserver<FlinkSessionJob> {
 
     private static final Logger LOG = 
LoggerFactory.getLogger(FlinkSessionJobObserver.class);
 
-    private final SessionJobStatusObserver jobStatusObserver;
+    private final JobStatusObserver<FlinkSessionJob> jobStatusObserver;
     private final SnapshotObserver<FlinkSessionJob, FlinkSessionJobStatus> 
savepointObserver;
 
     public FlinkSessionJobObserver(EventRecorder eventRecorder) {
         super(eventRecorder);
-        this.jobStatusObserver = new SessionJobStatusObserver(eventRecorder);
+        this.jobStatusObserver = new JobStatusObserver<>(eventRecorder);
         this.savepointObserver = new SnapshotObserver<>(eventRecorder);
     }
 
@@ -70,90 +60,21 @@ public class FlinkSessionJobObserver extends 
AbstractFlinkResourceObserver<Flink
     @Override
     protected boolean 
checkIfAlreadyUpgraded(FlinkResourceContext<FlinkSessionJob> ctx) {
         var flinkSessionJob = ctx.getResource();
-        Collection<JobStatusMessage> jobStatusMessages;
         try {
-            jobStatusMessages = 
ctx.getFlinkService().listJobs(ctx.getObserveConfig());
-        } catch (Exception e) {
-            throw new RuntimeException("Failed to list jobs", e);
-        }
-        var submittedJobId = 
flinkSessionJob.getStatus().getJobStatus().getJobId();
-        for (JobStatusMessage jobStatusMessage : jobStatusMessages) {
-            if 
(jobStatusMessage.getJobId().toHexString().equals(submittedJobId)) {
-                LOG.info("Job with id {} is already deployed.", 
submittedJobId);
-                return true;
+            var jobStatus = flinkSessionJob.getStatus().getJobStatus();
+            if (jobStatus.getJobId() == null) {
+                // No job was submitted
+                return false;
             }
-        }
-        return false;
-    }
-
-    private static class SessionJobStatusObserver extends 
JobStatusObserver<FlinkSessionJob> {
-
-        public SessionJobStatusObserver(EventRecorder eventRecorder) {
-            super(eventRecorder);
-        }
-
-        @Override
-        protected void onTimeout(FlinkResourceContext<FlinkSessionJob> ctx) {}
-
-        @Override
-        protected Optional<JobStatusMessage> filterTargetJob(
-                JobStatus status, List<JobStatusMessage> clusterJobStatuses) {
-            var jobId =
-                    Preconditions.checkNotNull(
-                            status.getJobId(), "The jobID to be observed 
should not be null");
-            var matchedList =
-                    clusterJobStatuses.stream()
-                            .filter(job -> 
job.getJobId().toHexString().equals(jobId))
-                            .collect(Collectors.toList());
-            Preconditions.checkArgument(
-                    matchedList.size() <= 1,
-                    String.format(
-                            "Expected one job for JobID: %s, but found %d",
-                            status.getJobId(), matchedList.size()));
-
-            if (matchedList.size() == 0) {
-                LOG.warn("No job found for JobID: {}", jobId);
-                return Optional.empty();
+            var jobId = JobID.fromHexString(jobStatus.getJobId());
+            if (ctx.getFlinkService().getJobStatus(ctx.getObserveConfig(), 
jobId).isPresent()) {
+                LOG.info("Job with id {} is already deployed.", jobId);
+                return true;
             } else {
-                return Optional.of(matchedList.get(0));
+                return false;
             }
-        }
-
-        @Override
-        protected void 
onTargetJobNotFound(FlinkResourceContext<FlinkSessionJob> ctx) {
-            ifHaDisabledMarkSessionJobMissing(ctx);
-        }
-
-        @Override
-        protected void onNoJobsFound(FlinkResourceContext<FlinkSessionJob> 
ctx) {
-            ifHaDisabledMarkSessionJobMissing(ctx);
-        }
-
-        /**
-         * When HA is disabled the session job will not recover on JM 
restarts. If the JM goes down
-         * / restarted the session job should be marked missing.
-         *
-         * @param ctx Flink session job context.
-         */
-        private void 
ifHaDisabledMarkSessionJobMissing(FlinkResourceContext<FlinkSessionJob> ctx) {
-            var sessionJob = ctx.getResource();
-            if 
(HighAvailabilityMode.isHighAvailabilityModeActivated(ctx.getObserveConfig())) {
-                return;
-            }
-            sessionJob
-                    .getStatus()
-                    .getJobStatus()
-                    
.setState(org.apache.flink.api.common.JobStatus.RECONCILING.name());
-            LOG.error(MISSING_SESSION_JOB_ERR);
-            ReconciliationUtils.updateForReconciliationError(
-                    ctx, new 
MissingSessionJobException(MISSING_SESSION_JOB_ERR));
-            eventRecorder.triggerEvent(
-                    sessionJob,
-                    EventRecorder.Type.Warning,
-                    EventRecorder.Reason.Missing,
-                    EventRecorder.Component.Job,
-                    MISSING_SESSION_JOB_ERR,
-                    ctx.getKubernetesClient());
+        } catch (Exception e) {
+            throw new RuntimeException("Failed to list jobs", e);
         }
     }
 }
diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconciler.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconciler.java
index 8621c50b..7992a382 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconciler.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconciler.java
@@ -203,12 +203,14 @@ public class ApplicationReconciler
         // https://issues.apache.org/jira/browse/FLINK-19358
         // https://issues.apache.org/jira/browse/FLINK-29109
 
-        if (deployConfig.get(PipelineOptionsInternal.PIPELINE_FIXED_JOB_ID) != 
null) {
-            // user managed, don't touch
+        var status = resource.getStatus();
+        var userJobId = 
deployConfig.get(PipelineOptionsInternal.PIPELINE_FIXED_JOB_ID);
+        if (userJobId != null) {
+            status.getJobStatus().setJobId(userJobId);
+            statusRecorder.patchAndCacheStatus(resource, client);
             return;
         }
 
-        var status = resource.getStatus();
         // Rotate job id when not last-state deployment
         if (status.getJobStatus().getJobId() == null || !lastStateDeploy) {
             String jobId = JobID.generate().toHexString();
diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconciler.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconciler.java
index b71264d6..de4d59e4 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconciler.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconciler.java
@@ -30,10 +30,9 @@ import 
org.apache.flink.kubernetes.operator.autoscaler.KubernetesJobAutoScalerCo
 import 
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions;
 import org.apache.flink.kubernetes.operator.controller.FlinkResourceContext;
 import 
org.apache.flink.kubernetes.operator.reconciler.deployment.AbstractJobReconciler;
+import org.apache.flink.kubernetes.operator.service.AbstractFlinkService;
 import org.apache.flink.kubernetes.operator.utils.EventRecorder;
 import org.apache.flink.kubernetes.operator.utils.StatusRecorder;
-import org.apache.flink.runtime.messages.FlinkJobNotFoundException;
-import 
org.apache.flink.runtime.messages.FlinkJobTerminatedWithoutCancellationException;
 
 import io.javaoperatorsdk.operator.api.reconciler.DeleteControl;
 import io.javaoperatorsdk.operator.processing.event.ResourceID;
@@ -43,8 +42,6 @@ import org.slf4j.LoggerFactory;
 import java.util.Optional;
 import java.util.concurrent.ExecutionException;
 
-import static org.apache.flink.util.ExceptionUtils.findThrowable;
-
 /** The reconciler for the {@link FlinkSessionJob}. */
 public class SessionJobReconciler
         extends AbstractJobReconciler<FlinkSessionJob, FlinkSessionJobSpec, 
FlinkSessionJobStatus> {
@@ -129,19 +126,10 @@ public class SessionJobReconciler
                                     : UpgradeMode.STATELESS;
                     cancelJob(ctx, upgradeMode);
                 } catch (ExecutionException e) {
-                    if (findThrowable(e, 
FlinkJobNotFoundException.class).isPresent()) {
-                        LOG.error("Job {} not found in the Flink cluster.", 
jobID, e);
+                    if (AbstractFlinkService.isJobMissingOrTerminated(e)) {
                         return DeleteControl.defaultDelete();
                     }
-
-                    if (findThrowable(e, 
FlinkJobTerminatedWithoutCancellationException.class)
-                            .isPresent()) {
-                        LOG.error("Job {} already terminated without 
cancellation.", jobID, e);
-                        return DeleteControl.defaultDelete();
-                    }
-
-                    final long delay =
-                            
ctx.getOperatorConfig().getProgressCheckInterval().toMillis();
+                    long delay = 
ctx.getOperatorConfig().getProgressCheckInterval().toMillis();
                     LOG.error(
                             "Failed to cancel job {}, will reschedule after {} 
milliseconds.",
                             jobID,
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 ae8979a5..48088435 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
@@ -54,6 +54,8 @@ import org.apache.flink.runtime.client.JobStatusMessage;
 import 
org.apache.flink.runtime.highavailability.nonha.standalone.StandaloneClientHAServices;
 import org.apache.flink.runtime.jobgraph.RestoreMode;
 import org.apache.flink.runtime.jobmaster.JobResult;
+import org.apache.flink.runtime.messages.FlinkJobNotFoundException;
+import 
org.apache.flink.runtime.messages.FlinkJobTerminatedWithoutCancellationException;
 import org.apache.flink.runtime.rest.FileUpload;
 import org.apache.flink.runtime.rest.RestClient;
 import org.apache.flink.runtime.rest.handler.async.AsynchronousOperationResult;
@@ -78,6 +80,7 @@ import 
org.apache.flink.runtime.rest.messages.job.savepoints.SavepointTriggerReq
 import org.apache.flink.runtime.rest.messages.queue.QueueStatus;
 import org.apache.flink.runtime.rest.messages.taskmanager.TaskManagersHeaders;
 import org.apache.flink.runtime.rest.messages.taskmanager.TaskManagersInfo;
+import org.apache.flink.runtime.rest.util.RestClientException;
 import org.apache.flink.runtime.rest.util.RestConstants;
 import 
org.apache.flink.runtime.scheduler.stopwithsavepoint.StopWithSavepointStoppingException;
 import 
org.apache.flink.runtime.state.memory.NonPersistentMetadataCheckpointStorageLocation;
@@ -95,6 +98,8 @@ import org.apache.flink.util.FlinkException;
 import org.apache.flink.util.FlinkRuntimeException;
 import org.apache.flink.util.Preconditions;
 
+import 
org.apache.flink.shaded.netty4.io.netty.handler.codec.http.HttpResponseStatus;
+
 import io.fabric8.kubernetes.api.model.DeletionPropagation;
 import io.fabric8.kubernetes.api.model.ObjectMeta;
 import io.fabric8.kubernetes.api.model.PodList;
@@ -123,7 +128,6 @@ import java.nio.file.Files;
 import java.nio.file.Paths;
 import java.time.Duration;
 import java.util.Arrays;
-import java.util.Collection;
 import java.util.Collections;
 import java.util.HashMap;
 import java.util.List;
@@ -141,6 +145,7 @@ import java.util.stream.Collectors;
 
 import static 
org.apache.flink.kubernetes.operator.config.FlinkConfigBuilder.FLINK_VERSION;
 import static 
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions.K8S_OP_CONF_PREFIX;
+import static org.apache.flink.util.ExceptionUtils.findThrowable;
 
 /**
  * An abstract {@link FlinkService} containing some common implementations for 
the native and
@@ -263,14 +268,24 @@ public abstract class AbstractFlinkService implements 
FlinkService {
     }
 
     @Override
-    public Collection<JobStatusMessage> listJobs(Configuration conf) throws 
Exception {
+    public Optional<JobStatusMessage> getJobStatus(Configuration conf, JobID 
jobId)
+            throws Exception {
         try (var clusterClient = getClusterClient(conf)) {
             return clusterClient
                     .sendRequest(
                             JobsOverviewHeaders.getInstance(),
                             EmptyMessageParameters.getInstance(),
                             EmptyRequestBody.getInstance())
-                    .thenApply(JobStatusUtils::toJobStatusMessage)
+                    .thenApply(
+                            mjd -> {
+                                if (mjd.getJobs() == null) {
+                                    return Optional.<JobStatusMessage>empty();
+                                }
+                                return mjd.getJobs().stream()
+                                        .filter(jd -> 
jd.getJobId().equals(jobId))
+                                        .findAny()
+                                        
.map(JobStatusUtils::toJobStatusMessage);
+                            })
                     .get(operatorConfig.getFlinkClientTimeout().toSeconds(), 
TimeUnit.SECONDS);
         }
     }
@@ -425,16 +440,23 @@ public abstract class AbstractFlinkService implements 
FlinkService {
             LOG.debug("Job is not in terminal state, cancelling it");
 
             try (var clusterClient = getClusterClient(conf)) {
-                final String clusterId = clusterClient.getClusterId();
                 switch (upgradeMode) {
                     case STATELESS:
                         LOG.info("Cancelling job.");
-                        clusterClient
-                                .cancel(jobId)
-                                .get(
-                                        
operatorConfig.getFlinkCancelJobTimeout().toSeconds(),
-                                        TimeUnit.SECONDS);
-                        LOG.info("Job successfully cancelled.");
+                        try {
+                            clusterClient
+                                    .cancel(jobId)
+                                    .get(
+                                            
operatorConfig.getFlinkCancelJobTimeout().toSeconds(),
+                                            TimeUnit.SECONDS);
+                            LOG.info("Job successfully cancelled.");
+                        } catch (Exception e) {
+                            if (isJobMissingOrTerminated(e)) {
+                                LOG.info("Job already missing or terminated");
+                            } else {
+                                throw e;
+                            }
+                        }
                         break;
                     case SAVEPOINT:
                         if 
(ReconciliationUtils.isJobRunning(sessionJobStatus)) {
@@ -464,10 +486,9 @@ public abstract class AbstractFlinkService implements 
FlinkService {
                             } catch (TimeoutException exception) {
                                 throw new FlinkException(
                                         String.format(
-                                                "Timed out stopping the job %s 
in Flink cluster %s with savepoint, "
+                                                "Timed out stopping the job %s 
with savepoint, "
                                                         + "please configure a 
larger timeout via '%s'",
                                                 jobId,
-                                                clusterId,
                                                 
ExecutionCheckpointingOptions.CHECKPOINTING_TIMEOUT
                                                         .key()),
                                         exception);
@@ -494,6 +515,22 @@ public abstract class AbstractFlinkService implements 
FlinkService {
                 });
     }
 
+    public static boolean isJobMissingOrTerminated(Exception e) {
+        if (findThrowable(e, FlinkJobNotFoundException.class).isPresent()
+                || findThrowable(e, 
FlinkJobTerminatedWithoutCancellationException.class)
+                        .isPresent()) {
+            return true;
+        }
+
+        return findThrowable(e, RestClientException.class)
+                .map(RestClientException::getHttpResponseStatus)
+                .map(
+                        respCode ->
+                                HttpResponseStatus.NOT_FOUND == respCode
+                                        || HttpResponseStatus.CONFLICT == 
respCode)
+                .orElse(false);
+    }
+
     @Override
     public void triggerSavepoint(
             String jobId,
diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/FlinkService.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/FlinkService.java
index 3c66cd1b..0e7bfca0 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/FlinkService.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/FlinkService.java
@@ -43,7 +43,6 @@ import io.fabric8.kubernetes.client.KubernetesClient;
 
 import javax.annotation.Nullable;
 
-import java.util.Collection;
 import java.util.List;
 import java.util.Map;
 import java.util.Optional;
@@ -70,7 +69,7 @@ public interface FlinkService {
 
     boolean isJobManagerPortReady(Configuration config);
 
-    Collection<JobStatusMessage> listJobs(Configuration conf) throws Exception;
+    Optional<JobStatusMessage> getJobStatus(Configuration conf, JobID jobId) 
throws Exception;
 
     JobResult requestJobResult(Configuration conf, JobID jobID) throws 
Exception;
 
diff --git 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/TestingFlinkService.java
 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/TestingFlinkService.java
index 6e461e4d..3283191d 100644
--- 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/TestingFlinkService.java
+++ 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/TestingFlinkService.java
@@ -276,11 +276,12 @@ public class TestingFlinkService extends 
AbstractFlinkService {
     }
 
     @Override
-    public Collection<JobStatusMessage> listJobs(Configuration conf) throws 
Exception {
+    public Optional<JobStatusMessage> getJobStatus(Configuration conf, JobID 
jobID)
+            throws Exception {
         if (!isPortReady) {
             throw new TimeoutException("JM port is unavailable");
         }
-        return super.listJobs(conf);
+        return super.getJobStatus(conf, jobID);
     }
 
     public List<Tuple3<String, JobStatusMessage, Configuration>> listJobs() {
@@ -603,7 +604,10 @@ public class TestingFlinkService extends 
AbstractFlinkService {
     }
 
     @Override
-    public Map<String, String> getClusterInfo(Configuration conf) {
+    public Map<String, String> getClusterInfo(Configuration conf) throws 
TimeoutException {
+        if (!isPortReady) {
+            throw new TimeoutException("JM port is unavailable");
+        }
         return CLUSTER_INFO;
     }
 
diff --git 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/observer/deployment/ApplicationObserverTest.java
 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/observer/deployment/ApplicationObserverTest.java
index 03062a8c..8e7e2c24 100644
--- 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/observer/deployment/ApplicationObserverTest.java
+++ 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/observer/deployment/ApplicationObserverTest.java
@@ -20,6 +20,7 @@ package 
org.apache.flink.kubernetes.operator.observer.deployment;
 import org.apache.flink.api.common.JobID;
 import org.apache.flink.configuration.Configuration;
 import org.apache.flink.configuration.PipelineOptions;
+import org.apache.flink.configuration.PipelineOptionsInternal;
 import org.apache.flink.kubernetes.operator.OperatorTestBase;
 import org.apache.flink.kubernetes.operator.TestUtils;
 import org.apache.flink.kubernetes.operator.TestingFlinkService;
@@ -72,16 +73,23 @@ public class ApplicationObserverTest extends 
OperatorTestBase {
 
     private Context<FlinkDeployment> readyContext;
     private TestObserverAdapter<FlinkDeployment> observer;
+    private FlinkDeployment deployment;
 
     @Override
     public void setup() {
         observer = new TestObserverAdapter<>(this, new 
ApplicationObserver(eventRecorder));
         readyContext = 
TestUtils.createContextWithReadyJobManagerDeployment(kubernetesClient);
+        deployment = TestUtils.buildApplicationCluster();
+        var jobId = new JobID().toHexString();
+        deployment
+                .getSpec()
+                .getFlinkConfiguration()
+                .put(PipelineOptionsInternal.PIPELINE_FIXED_JOB_ID.key(), 
jobId);
+        deployment.getStatus().getJobStatus().setJobId(jobId);
     }
 
     @Test
     public void observeApplicationCluster() throws Exception {
-        FlinkDeployment deployment = TestUtils.buildApplicationCluster();
         Configuration conf =
                 configManager.getDeployConfig(deployment.getMetadata(), 
deployment.getSpec());
 
@@ -178,12 +186,10 @@ public class ApplicationObserverTest extends 
OperatorTestBase {
 
     @Test
     public void testEventGeneratedWhenStatusChanged() throws Exception {
-        FlinkDeployment deployment = TestUtils.buildApplicationCluster();
         Configuration conf =
                 configManager.getDeployConfig(deployment.getMetadata(), 
deployment.getSpec());
         flinkService.submitApplicationCluster(deployment.getSpec().getJob(), 
conf, false);
 
-        deployment.setStatus(deployment.initStatus());
         ReconciliationUtils.updateStatusForDeployedSpec(deployment, new 
Configuration());
         
deployment.getStatus().setJobManagerDeploymentStatus(JobManagerDeploymentStatus.READY);
 
@@ -211,12 +217,10 @@ public class ApplicationObserverTest extends 
OperatorTestBase {
 
     @Test
     public void testErrorForwardToStatusWhenJobFailed() throws Exception {
-        FlinkDeployment deployment = TestUtils.buildApplicationCluster();
         Configuration conf =
                 configManager.getDeployConfig(deployment.getMetadata(), 
deployment.getSpec());
         flinkService.submitApplicationCluster(deployment.getSpec().getJob(), 
conf, false);
 
-        deployment.setStatus(deployment.initStatus());
         ReconciliationUtils.updateStatusForDeployedSpec(deployment, new 
Configuration());
         
deployment.getStatus().setJobManagerDeploymentStatus(JobManagerDeploymentStatus.READY);
 
@@ -232,7 +236,6 @@ public class ApplicationObserverTest extends 
OperatorTestBase {
 
     @Test
     public void observeSavepoint() throws Exception {
-        FlinkDeployment deployment = TestUtils.buildApplicationCluster();
         Long timedOutNonce = 1L;
         deployment.getSpec().getJob().setSavepointTriggerNonce(timedOutNonce);
         Configuration conf =
@@ -524,7 +527,6 @@ public class ApplicationObserverTest extends 
OperatorTestBase {
 
     @Test
     public void observeCheckpoint() throws Exception {
-        FlinkDeployment deployment = TestUtils.buildApplicationCluster();
         Long timedOutNonce = 1L;
         final var errorReason = EventRecorder.Reason.CheckpointError.name();
         final var namespace = deployment.getMetadata().getNamespace();
@@ -656,7 +658,6 @@ public class ApplicationObserverTest extends 
OperatorTestBase {
 
     @Test
     public void testSavepointFormat() throws Exception {
-        FlinkDeployment deployment = TestUtils.buildApplicationCluster();
         Configuration conf =
                 configManager.getDeployConfig(deployment.getMetadata(), 
deployment.getSpec());
         flinkService.submitApplicationCluster(deployment.getSpec().getJob(), 
conf, false);
@@ -731,9 +732,8 @@ public class ApplicationObserverTest extends 
OperatorTestBase {
 
     private void bringToReadyStatus(FlinkDeployment deployment) {
         ReconciliationUtils.updateStatusForDeployedSpec(deployment, new 
Configuration());
-        JobStatus jobStatus = new JobStatus();
+        JobStatus jobStatus = deployment.getStatus().getJobStatus();
         jobStatus.setJobName("jobname");
-        jobStatus.setJobId("0000000000");
         jobStatus.setState(JobState.RUNNING.name());
         deployment.getStatus().setJobStatus(jobStatus);
         
deployment.getStatus().setJobManagerDeploymentStatus(JobManagerDeploymentStatus.READY);
@@ -741,7 +741,6 @@ public class ApplicationObserverTest extends 
OperatorTestBase {
 
     @Test
     public void observeListJobsError() {
-        FlinkDeployment deployment = TestUtils.buildApplicationCluster();
         bringToReadyStatus(deployment);
         observer.observe(deployment, readyContext);
         assertEquals(
@@ -771,7 +770,6 @@ public class ApplicationObserverTest extends 
OperatorTestBase {
 
         var context = 
TestUtils.createContextWithDeployment(kubernetesDeployment, kubernetesClient);
 
-        FlinkDeployment deployment = TestUtils.buildApplicationCluster();
         deployment.getMetadata().setGeneration(123L);
 
         var status = deployment.getStatus();
@@ -849,8 +847,6 @@ public class ApplicationObserverTest extends 
OperatorTestBase {
 
     @Test
     public void observeAlreadyScaled() {
-        var deployment = TestUtils.buildApplicationCluster();
-
         // Update status for for running job
         ReconciliationUtils.updateStatusBeforeDeploymentAttempt(
                 deployment,
@@ -880,7 +876,6 @@ public class ApplicationObserverTest extends 
OperatorTestBase {
 
     @Test
     public void validateLastReconciledClearedOnInitialFailure() {
-        FlinkDeployment deployment = TestUtils.buildApplicationCluster();
         deployment.getMetadata().setGeneration(123L);
 
         ReconciliationUtils.updateStatusBeforeDeploymentAttempt(
diff --git 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/observer/sessionjob/FlinkSessionJobObserverTest.java
 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/observer/sessionjob/FlinkSessionJobObserverTest.java
index c0005127..3898d996 100644
--- 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/observer/sessionjob/FlinkSessionJobObserverTest.java
+++ 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/observer/sessionjob/FlinkSessionJobObserverTest.java
@@ -21,7 +21,6 @@ import org.apache.flink.api.common.JobID;
 import org.apache.flink.api.common.JobStatus;
 import org.apache.flink.autoscaler.NoopJobAutoscaler;
 import org.apache.flink.configuration.Configuration;
-import org.apache.flink.configuration.HighAvailabilityOptions;
 import org.apache.flink.configuration.RestOptions;
 import org.apache.flink.kubernetes.operator.OperatorTestBase;
 import org.apache.flink.kubernetes.operator.TestUtils;
@@ -40,7 +39,6 @@ import 
org.apache.flink.kubernetes.operator.utils.SnapshotUtils;
 import io.fabric8.kubernetes.client.KubernetesClient;
 import io.fabric8.kubernetes.client.server.mock.EnableKubernetesMockClient;
 import lombok.Getter;
-import org.apache.commons.lang3.StringUtils;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.params.ParameterizedTest;
@@ -139,29 +137,13 @@ public class FlinkSessionJobObserverTest extends 
OperatorTestBase {
 
         eventCollector.events.clear();
 
-        // With HA enabled no error should be triggered
-        observer.observe(sessionJob2, readyContext);
-        Assertions.assertEquals(
-                JobStatus.RECONCILING.name(), 
sessionJob2.getStatus().getJobStatus().getState());
-        
Assertions.assertTrue(StringUtils.isEmpty(sessionJob2.getStatus().getError()));
-        Assertions.assertTrue(eventCollector.events.isEmpty());
-
-        // With HA disabled we expect an error status and event
-        sessionJob2
-                .getSpec()
-                .getFlinkConfiguration()
-                .put(HighAvailabilityOptions.HA_MODE.key(), "NONE");
         observer.observe(sessionJob2, readyContext);
         Assertions.assertEquals(
                 JobStatus.RECONCILING.name(), 
sessionJob2.getStatus().getJobStatus().getState());
         Assertions.assertTrue(
-                sessionJob2
-                        .getStatus()
-                        .getError()
-                        .contains(JobStatusObserver.MISSING_SESSION_JOB_ERR));
+                
sessionJob2.getStatus().getError().contains(JobStatusObserver.JOB_NOT_FOUND_ERR));
         Assertions.assertEquals(
-                JobStatusObserver.MISSING_SESSION_JOB_ERR,
-                eventCollector.events.peek().getMessage());
+                JobStatusObserver.JOB_NOT_FOUND_ERR, 
eventCollector.events.peek().getMessage());
     }
 
     @Test
@@ -179,7 +161,7 @@ public class FlinkSessionJobObserverTest extends 
OperatorTestBase {
                 JobStatus.RECONCILING.name(), 
sessionJob.getStatus().getJobStatus().getState());
 
         flinkService.setListJobConsumer(
-                configuration ->
+                (configuration) ->
                         Assertions.assertEquals(8088, 
configuration.getInteger(RestOptions.PORT)));
         observer.observe(sessionJob, readyContext);
         Assertions.assertEquals(
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 c4f24ee9..a6ff6e4b 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
@@ -77,6 +77,7 @@ import 
org.apache.flink.runtime.rest.messages.job.savepoints.SavepointTriggerReq
 import org.apache.flink.runtime.rest.messages.taskmanager.TaskManagerInfo;
 import org.apache.flink.runtime.rest.messages.taskmanager.TaskManagersHeaders;
 import org.apache.flink.runtime.rest.messages.taskmanager.TaskManagersInfo;
+import org.apache.flink.runtime.rest.util.RestClientException;
 import org.apache.flink.runtime.rest.util.RestMapperUtils;
 import 
org.apache.flink.runtime.scheduler.stopwithsavepoint.StopWithSavepointStoppingException;
 import org.apache.flink.runtime.taskexecutor.TaskExecutorMemoryConfiguration;
@@ -91,6 +92,7 @@ import org.apache.flink.util.concurrent.Executors;
 import org.apache.flink.util.function.TriFunction;
 
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper;
+import 
org.apache.flink.shaded.netty4.io.netty.handler.codec.http.HttpResponseStatus;
 
 import io.fabric8.kubernetes.api.model.DeletionPropagation;
 import io.fabric8.kubernetes.api.model.ListMeta;
@@ -301,6 +303,40 @@ public class AbstractFlinkServiceTest {
         assertNull(jobStatus.getSavepointInfo().getLastSavepoint());
     }
 
+    @ParameterizedTest
+    @ValueSource(ints = {404, 409, 500})
+    public void cancelErrorHandling(int statusCode) throws Exception {
+
+        var testingClusterClient =
+                new TestingClusterClient<>(configuration, 
TestUtils.TEST_DEPLOYMENT_NAME);
+        testingClusterClient.setCancelFunction(
+                jobID ->
+                        CompletableFuture.failedFuture(
+                                new RuntimeException(
+                                        new RestClientException(
+                                                "errrr", 
HttpResponseStatus.valueOf(statusCode)))));
+        var flinkService = new TestingService(testingClusterClient);
+
+        JobID jobID = JobID.generate();
+        var job = TestUtils.buildSessionJob();
+        var jobStatus = job.getStatus().getJobStatus();
+        jobStatus.setJobId(jobID.toHexString());
+        jobStatus.setState("RUNNING");
+        ReconciliationUtils.updateStatusForDeployedSpec(job, new 
Configuration());
+
+        if (statusCode == 500) {
+            assertThrows(
+                    Exception.class,
+                    () ->
+                            flinkService.cancelSessionJob(
+                                    job, UpgradeMode.STATELESS, new 
Configuration()));
+            assertEquals("RUNNING", jobStatus.getState());
+        } else {
+            flinkService.cancelSessionJob(job, UpgradeMode.STATELESS, new 
Configuration());
+            assertEquals("FINISHED", jobStatus.getState());
+        }
+    }
+
     @ParameterizedTest
     @ValueSource(booleans = {true, false})
     public void cancelJobWithSavepointUpgradeModeTest(boolean 
deleteAfterSavepoint)

Reply via email to