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)