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 207b149f [FLINK-35292] Set dummy savepoint path during last-state
upgrade
207b149f is described below
commit 207b149f9b556d68c4ab98d16cdde2f7820659c0
Author: Gyula Fora <[email protected]>
AuthorDate: Mon May 6 10:14:29 2024 +0200
[FLINK-35292] Set dummy savepoint path during last-state upgrade
---
.../deployment/AbstractJobReconciler.java | 11 ++++-
.../deployment/ApplicationReconciler.java | 12 +++++
.../operator/service/AbstractFlinkService.java | 9 ++++
.../kubernetes/operator/service/FlinkService.java | 2 +
.../kubernetes/operator/utils/FlinkUtils.java | 24 ++++++++-
.../kubernetes/operator/utils/SnapshotUtils.java | 21 ++++++++
.../kubernetes/operator/TestingFlinkService.java | 6 +++
.../deployment/ApplicationReconcilerTest.java | 23 ++++++++-
.../ApplicationReconcilerUpgradeModeTest.java | 54 +++++++++++++++++++-
.../kubernetes/operator/utils/FlinkUtilsTest.java | 57 ++++++++++++++++++++++
10 files changed, 214 insertions(+), 5 deletions(-)
diff --git
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/AbstractJobReconciler.java
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/AbstractJobReconciler.java
index 1e256b7a..bdaec55f 100644
---
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/AbstractJobReconciler.java
+++
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/AbstractJobReconciler.java
@@ -33,6 +33,7 @@ import
org.apache.flink.kubernetes.operator.api.status.SnapshotTriggerType;
import
org.apache.flink.kubernetes.operator.autoscaler.KubernetesJobAutoScalerContext;
import
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions;
import org.apache.flink.kubernetes.operator.controller.FlinkResourceContext;
+import org.apache.flink.kubernetes.operator.exception.RecoveryFailureException;
import org.apache.flink.kubernetes.operator.reconciler.ReconciliationUtils;
import org.apache.flink.kubernetes.operator.reconciler.SnapshotType;
import org.apache.flink.kubernetes.operator.service.CheckpointHistoryWrapper;
@@ -64,6 +65,8 @@ public abstract class AbstractJobReconciler<
private static final Logger LOG =
LoggerFactory.getLogger(AbstractJobReconciler.class);
+ public static final String LAST_STATE_DUMMY_SP_PATH =
"KUBERNETES_OPERATOR_LAST_STATE";
+
public AbstractJobReconciler(
EventRecorder eventRecorder,
StatusRecorder<CR, STATUS> statusRecorder,
@@ -179,6 +182,12 @@ public abstract class AbstractJobReconciler<
var flinkService = ctx.getFlinkService();
if (ReconciliationUtils.isJobInTerminalState(status)
&&
!flinkService.isHaMetadataAvailable(ctx.getObserveConfig())) {
+
+ if (!SnapshotUtils.lastSavepointKnown(status)) {
+ throw new RecoveryFailureException(
+ "Job is in terminal state but last checkpoint is
unknown, possibly due to an unrecoverable restore error. Manual restore
required.",
+ "UpgradeFailed");
+ }
LOG.info(
"Job is in terminal state, ready for upgrade from observed
latest checkpoint/savepoint");
return AvailableUpgradeMode.of(UpgradeMode.SAVEPOINT);
@@ -265,7 +274,7 @@ public abstract class AbstractJobReconciler<
throws Exception {
Optional<String> savepointOpt = Optional.empty();
- if (spec.getJob().getUpgradeMode() != UpgradeMode.STATELESS) {
+ if (spec.getJob().getUpgradeMode() == UpgradeMode.SAVEPOINT) {
savepointOpt =
Optional.ofNullable(
ctx.getResource()
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 7992a382..72f3a51f 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
@@ -29,6 +29,8 @@ import
org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec;
import org.apache.flink.kubernetes.operator.api.spec.UpgradeMode;
import org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentStatus;
import
org.apache.flink.kubernetes.operator.api.status.JobManagerDeploymentStatus;
+import org.apache.flink.kubernetes.operator.api.status.Savepoint;
+import org.apache.flink.kubernetes.operator.api.status.SnapshotTriggerType;
import
org.apache.flink.kubernetes.operator.autoscaler.KubernetesJobAutoScalerContext;
import
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions;
import org.apache.flink.kubernetes.operator.controller.FlinkResourceContext;
@@ -159,8 +161,18 @@ public class ApplicationReconciler
relatedResource.getStatus().getClusterInfo());
if (savepoint.isPresent()) {
+ // Savepoint deployment
deployConfig.set(SavepointConfigOptions.SAVEPOINT_PATH,
savepoint.get());
+ } else if (requireHaMetadata &&
flinkService.atLeastOneCheckpoint(deployConfig)) {
+ // Last state deployment, explicitly set a dummy savepoint path to
avoid accidental
+ // incorrect state restore in case the HA metadata is deleted by
the user
+ deployConfig.set(SavepointConfigOptions.SAVEPOINT_PATH,
LAST_STATE_DUMMY_SP_PATH);
+ status.getJobStatus()
+ .getSavepointInfo()
+ .setLastSavepoint(
+ Savepoint.of(LAST_STATE_DUMMY_SP_PATH,
SnapshotTriggerType.UNKNOWN));
} else {
+ // Stateless deployment, remove any user configured savepoint path
deployConfig.removeConfig(SavepointConfigOptions.SAVEPOINT_PATH);
}
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 48088435..001b9673 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
@@ -229,6 +229,15 @@ public abstract class AbstractFlinkService implements
FlinkService {
return false;
}
+ @Override
+ public boolean atLeastOneCheckpoint(Configuration conf) {
+ if (FlinkUtils.isKubernetesHAActivated(conf)) {
+ return
FlinkUtils.isKubernetesHaMetadataAvailableWithCheckpoint(conf,
kubernetesClient);
+ } else {
+ return isHaMetadataAvailable(conf);
+ }
+ }
+
@Override
public JobID submitJobToSessionCluster(
ObjectMeta meta,
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 0e7bfca0..1e2630d6 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
@@ -57,6 +57,8 @@ public interface FlinkService {
boolean isHaMetadataAvailable(Configuration conf);
+ boolean atLeastOneCheckpoint(Configuration conf);
+
void submitSessionCluster(Configuration conf) throws Exception;
JobID submitJobToSessionCluster(
diff --git
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/FlinkUtils.java
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/FlinkUtils.java
index 10288c70..3d81e509 100644
---
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/FlinkUtils.java
+++
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/FlinkUtils.java
@@ -58,6 +58,7 @@ import java.util.Iterator;
import java.util.LinkedHashMap;
import java.util.Map;
import java.util.Optional;
+import java.util.function.Predicate;
import static
org.apache.flink.kubernetes.utils.Constants.LABEL_CONFIGMAP_TYPE_HIGH_AVAILABILITY;
@@ -287,6 +288,20 @@ public class FlinkUtils {
public static boolean isKubernetesHaMetadataAvailable(
Configuration conf, KubernetesClient kubernetesClient) {
+ return isKubernetesHaMetadataAvailable(
+ conf, kubernetesClient, FlinkUtils::isValidHaConfigMap);
+ }
+
+ public static boolean isKubernetesHaMetadataAvailableWithCheckpoint(
+ Configuration conf, KubernetesClient kubernetesClient) {
+ return isKubernetesHaMetadataAvailable(
+ conf, kubernetesClient, cm -> isValidHaConfigMap(cm) &&
checkpointExists(cm));
+ }
+
+ private static boolean isKubernetesHaMetadataAvailable(
+ Configuration conf,
+ KubernetesClient kubernetesClient,
+ Predicate<ConfigMap> cmPredicate) {
String clusterId = conf.get(KubernetesConfigOptions.CLUSTER_ID);
String namespace = conf.get(KubernetesConfigOptions.NAMESPACE);
@@ -303,7 +318,7 @@ public class FlinkUtils {
.list()
.getItems();
- return configMaps.stream().anyMatch(FlinkUtils::isValidHaConfigMap);
+ return configMaps.stream().anyMatch(cmPredicate);
}
private static boolean isValidHaConfigMap(ConfigMap cm) {
@@ -319,6 +334,13 @@ public class FlinkUtils {
return name.endsWith("-jobmanager-leader");
}
+ private static boolean checkpointExists(ConfigMap cm) {
+ var data = cm.getData();
+ return data != null
+ && data.keySet().stream()
+ .anyMatch(s ->
s.startsWith(Constants.CHECKPOINT_ID_KEY_PREFIX));
+ }
+
private static boolean isJobGraphKey(Map.Entry<String, String> entry) {
return entry.getKey().startsWith(Constants.JOB_GRAPH_STORE_KEY_PREFIX);
}
diff --git
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/SnapshotUtils.java
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/SnapshotUtils.java
index 2322e484..ea9345ed 100644
---
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/SnapshotUtils.java
+++
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/SnapshotUtils.java
@@ -22,12 +22,14 @@ import org.apache.flink.configuration.Configuration;
import org.apache.flink.configuration.ConfigurationUtils;
import org.apache.flink.kubernetes.operator.api.AbstractFlinkResource;
import org.apache.flink.kubernetes.operator.api.spec.FlinkVersion;
+import org.apache.flink.kubernetes.operator.api.status.CommonStatus;
import org.apache.flink.kubernetes.operator.api.status.JobStatus;
import org.apache.flink.kubernetes.operator.api.status.SnapshotInfo;
import org.apache.flink.kubernetes.operator.api.status.SnapshotTriggerType;
import
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions;
import org.apache.flink.kubernetes.operator.reconciler.ReconciliationUtils;
import org.apache.flink.kubernetes.operator.reconciler.SnapshotType;
+import
org.apache.flink.kubernetes.operator.reconciler.deployment.AbstractJobReconciler;
import org.apache.flink.kubernetes.operator.service.FlinkService;
import io.fabric8.kubernetes.client.KubernetesClient;
@@ -402,4 +404,23 @@ public class SnapshotUtils {
}
}
}
+
+ /**
+ * Check if the last snapshot information is known. True if the snapshot
location is known
+ * explicitly (not implicitly through a last-state upgrade) or if the
savepoint is known to be
+ * empty.
+ *
+ * @param status Flink resource status
+ * @return True if last savepoint is known
+ */
+ public static boolean lastSavepointKnown(CommonStatus<?> status) {
+ var lastSavepoint =
status.getJobStatus().getSavepointInfo().getLastSavepoint();
+
+ if (lastSavepoint == null) {
+ return true;
+ }
+ String location = lastSavepoint.getLocation();
+
+ return
!location.equals(AbstractJobReconciler.LAST_STATE_DUMMY_SP_PATH);
+ }
}
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 3283191d..b7427b01 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
@@ -125,6 +125,7 @@ public class TestingFlinkService extends
AbstractFlinkService {
@Setter private boolean isFlinkJobTerminatedWithoutCancellation = false;
@Setter private boolean isPortReady = true;
@Setter private boolean haDataAvailable = true;
+ @Setter private boolean checkpointAvailable = true;
@Setter private boolean jobManagerReady = true;
@Setter private boolean deployFailure = false;
@Setter private Runnable sessionJobSubmittedCallback;
@@ -236,6 +237,11 @@ public class TestingFlinkService extends
AbstractFlinkService {
}
}
+ @Override
+ public boolean atLeastOneCheckpoint(Configuration conf) {
+ return isHaMetadataAvailable(conf) && checkpointAvailable;
+ }
+
@Override
public boolean isHaMetadataAvailable(Configuration conf) {
return HighAvailabilityMode.isHighAvailabilityModeActivated(conf) &&
haDataAvailable;
diff --git
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerTest.java
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerTest.java
index 6f9e17fd..2840dc32 100644
---
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerTest.java
+++
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerTest.java
@@ -635,8 +635,15 @@ public class ApplicationReconcilerTest extends
OperatorTestBase {
private void verifyAndSetRunningJobsToStatus(
FlinkDeployment deployment,
List<Tuple3<String, JobStatusMessage, Configuration>> runningJobs)
{
+ verifyAndSetRunningJobsToStatus(deployment, runningJobs, null);
+ }
+
+ private void verifyAndSetRunningJobsToStatus(
+ FlinkDeployment deployment,
+ List<Tuple3<String, JobStatusMessage, Configuration>> runningJobs,
+ String savepoint) {
assertEquals(1, runningJobs.size());
- assertNull(runningJobs.get(0).f0);
+ assertEquals(savepoint, runningJobs.get(0).f0);
deployment
.getStatus()
.setJobStatus(
@@ -1124,7 +1131,10 @@ public class ApplicationReconcilerTest extends
OperatorTestBase {
deployment.getSpec().getRestartNonce(),
lastReconciledSpec.getRestartNonce());
// Set to running to let savepoint upgrade proceed
- verifyAndSetRunningJobsToStatus(deployment, flinkService.listJobs());
+ verifyAndSetRunningJobsToStatus(
+ deployment,
+ flinkService.listJobs(),
+ ApplicationReconciler.LAST_STATE_DUMMY_SP_PATH);
reconciler.reconcile(deployment, context);
// Make sure upgrade is properly triggered now
@@ -1132,6 +1142,15 @@ public class ApplicationReconcilerTest extends
OperatorTestBase {
deployment.getStatus().getReconciliationStatus().deserializeLastReconciledSpec();
assertEquals(deployment.getSpec().getRestartNonce(),
lastReconciledSpec.getRestartNonce());
assertEquals(JobState.SUSPENDED,
lastReconciledSpec.getJob().getState());
+ assertEquals(UpgradeMode.SAVEPOINT,
lastReconciledSpec.getJob().getUpgradeMode());
+ assertEquals(
+ "savepoint_0",
+ deployment
+ .getStatus()
+ .getJobStatus()
+ .getSavepointInfo()
+ .getLastSavepoint()
+ .getLocation());
}
@Test
diff --git
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerUpgradeModeTest.java
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerUpgradeModeTest.java
index a5416a04..ae32f197 100644
---
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerUpgradeModeTest.java
+++
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerUpgradeModeTest.java
@@ -52,6 +52,7 @@ import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.Arguments;
import org.junit.jupiter.params.provider.MethodSource;
+import org.junit.jupiter.params.provider.ValueSource;
import java.time.Duration;
import java.util.ArrayList;
@@ -568,6 +569,48 @@ public class ApplicationReconcilerUpgradeModeTest extends
OperatorTestBase {
assertEquals(JobState.SUSPENDED,
lastReconciledSpec.getJob().getState());
}
+ @ParameterizedTest
+ @ValueSource(booleans = {true, false})
+ public void testLastStateDummySpPath(boolean checkpointAvailable) throws
Exception {
+ // Bootstrap running deployment
+ var deployment = TestUtils.buildApplicationCluster();
+ deployment.getSpec().getJob().setUpgradeMode(UpgradeMode.LAST_STATE);
+
+ reconciler.reconcile(deployment, context);
+ verifyAndSetRunningJobsToStatus(deployment, flinkService.listJobs());
+
+ flinkService.setHaDataAvailable(true);
+ flinkService.setCheckpointAvailable(checkpointAvailable);
+
+ // Submit upgrade
+ deployment.getSpec().setRestartNonce(123L);
+ reconciler.reconcile(deployment, context);
+ reconciler.reconcile(deployment, context);
+
+ var lastReconciledSpec =
+
deployment.getStatus().getReconciliationStatus().deserializeLastReconciledSpec();
+
+ // Make sure we correctly record upgrade mode to last state
+ assertEquals(UpgradeMode.LAST_STATE,
lastReconciledSpec.getJob().getUpgradeMode());
+
+ if (checkpointAvailable) {
+ assertEquals(
+ ApplicationReconciler.LAST_STATE_DUMMY_SP_PATH,
+ deployment
+ .getStatus()
+ .getJobStatus()
+ .getSavepointInfo()
+ .getLastSavepoint()
+ .getLocation());
+ assertEquals(
+ ApplicationReconciler.LAST_STATE_DUMMY_SP_PATH,
+ flinkService.listJobs().get(0).f0);
+ } else {
+
assertNull(deployment.getStatus().getJobStatus().getSavepointInfo().getLastSavepoint());
+ assertNull(flinkService.listJobs().get(0).f0);
+ }
+ }
+
@Test
public void testUpgradeModeChangeFromSavepointToLastState() throws
Exception {
final String expectedSavepointPath = "savepoint_0";
@@ -682,7 +725,16 @@ public class ApplicationReconcilerUpgradeModeTest extends
OperatorTestBase {
.getImage());
// Upgrade mode changes from stateless to last-state while HA enabled
previously should not
// trigger a savepoint
- assertNull(flinkService.listJobs().get(0).f0);
+ assertEquals(
+ ApplicationReconciler.LAST_STATE_DUMMY_SP_PATH,
+ deployment
+ .getStatus()
+ .getJobStatus()
+ .getSavepointInfo()
+ .getLastSavepoint()
+ .getLocation());
+ assertEquals(
+ ApplicationReconciler.LAST_STATE_DUMMY_SP_PATH,
flinkService.listJobs().get(0).f0);
}
public static FlinkDeployment buildApplicationCluster(
diff --git
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/FlinkUtilsTest.java
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/FlinkUtilsTest.java
index a4c97f6d..34fdc11b 100644
---
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/FlinkUtilsTest.java
+++
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/FlinkUtilsTest.java
@@ -201,6 +201,63 @@ public class FlinkUtilsTest {
kubernetesClient));
}
+ @Test
+ public void kubernetesHaMetaDataCheckpointCheckTest() {
+ var cr = TestUtils.buildApplicationCluster();
+ var confManager = new FlinkConfigManager(new Configuration());
+ assertFalse(
+ FlinkUtils.isKubernetesHaMetadataAvailableWithCheckpoint(
+ confManager.getDeployConfig(cr.getMetadata(),
cr.getSpec()),
+ kubernetesClient));
+
+ var withCheckpoint = Map.of("checkpointID-2", "p");
+ var withoutCheckpoint = Map.of("counter", "2");
+
+ // Wrong CM name
+ createHAConfigMapWithData(
+ cr.getMetadata().getName() + "-wrong-name",
+ cr.getMetadata().getNamespace(),
+ cr.getMetadata().getName(),
+ withCheckpoint);
+ assertFalse(
+ FlinkUtils.isKubernetesHaMetadataAvailableWithCheckpoint(
+ confManager.getDeployConfig(cr.getMetadata(),
cr.getSpec()),
+ kubernetesClient));
+
+ // Missing data
+ createHAConfigMapWithData(
+ cr.getMetadata().getName() + "-000000000000-config-map",
+ cr.getMetadata().getNamespace(),
+ cr.getMetadata().getName(),
+ null);
+ assertFalse(
+ FlinkUtils.isKubernetesHaMetadataAvailableWithCheckpoint(
+ confManager.getDeployConfig(cr.getMetadata(),
cr.getSpec()),
+ kubernetesClient));
+
+ // CM data without CP
+ createHAConfigMapWithData(
+ cr.getMetadata().getName() + "-000000000000-config-map",
+ cr.getMetadata().getNamespace(),
+ cr.getMetadata().getName(),
+ withoutCheckpoint);
+ assertFalse(
+ FlinkUtils.isKubernetesHaMetadataAvailableWithCheckpoint(
+ confManager.getDeployConfig(cr.getMetadata(),
cr.getSpec()),
+ kubernetesClient));
+
+ // CM data with CP
+ createHAConfigMapWithData(
+ cr.getMetadata().getName() + "-000000000000-config-map",
+ cr.getMetadata().getNamespace(),
+ cr.getMetadata().getName(),
+ withCheckpoint);
+ assertTrue(
+ FlinkUtils.isKubernetesHaMetadataAvailableWithCheckpoint(
+ confManager.getDeployConfig(cr.getMetadata(),
cr.getSpec()),
+ kubernetesClient));
+ }
+
@Test
public void testJmNeverStartedDetection() {
var jmDeployment = new Deployment();