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 5296d63c [FLINK-35857] Fix redeploy failed deployment without latest
checkpoint
5296d63c is described below
commit 5296d63c6a948f598530f10304f7846fa5ee6a6a
Author: chenyuzhi459 <[email protected]>
AuthorDate: Wed Jul 24 17:58:10 2024 +0800
[FLINK-35857] Fix redeploy failed deployment without latest checkpoint
---
.../deployment/AbstractJobReconciler.java | 10 +++
.../kubernetes/operator/TestingFlinkService.java | 10 +++
.../controller/FailedDeploymentRestartTest.java | 87 ++++++++++++++++++++--
3 files changed, 102 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 bdaec55f..237c8f04 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
@@ -322,8 +322,18 @@ public abstract class AbstractJobReconciler<
throws Exception {
LOG.info("Resubmitting Flink job...");
SPEC specToRecover =
ReconciliationUtils.getDeployedSpec(ctx.getResource());
+ Optional<Savepoint> lastSavepoint =
+ Optional.ofNullable(
+ ctx.getResource()
+ .getStatus()
+ .getJobStatus()
+ .getSavepointInfo()
+ .getLastSavepoint());
if (requireHaMetadata) {
specToRecover.getJob().setUpgradeMode(UpgradeMode.LAST_STATE);
+ } else if (ctx.getResource().getSpec().getJob().getUpgradeMode() !=
UpgradeMode.STATELESS
+ && lastSavepoint.isPresent()) {
+ specToRecover.getJob().setUpgradeMode(UpgradeMode.SAVEPOINT);
}
restoreJob(ctx, specToRecover, ctx.getObserveConfig(),
requireHaMetadata);
}
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 b7427b01..0fa9e8f7 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
@@ -348,6 +348,16 @@ public class TestingFlinkService extends
AbstractFlinkService {
if (checkpointTriggers.containsKey(triggerId)) {
if (checkpointTriggers.get(triggerId)) {
+ // Mark completed checkpoint
+ checkpointInfo =
+ Tuple2.of(
+ Optional.of(
+ new
CheckpointHistoryWrapper.CompletedCheckpointInfo(
+ checkpointCounter,
+ "ck_" + checkpointCounter,
+ System.currentTimeMillis())),
+ Optional.empty());
+
checkpointCounter++;
return CheckpointFetchResult.completed();
}
diff --git
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FailedDeploymentRestartTest.java
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FailedDeploymentRestartTest.java
index 131d0756..de5c1695 100644
---
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FailedDeploymentRestartTest.java
+++
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FailedDeploymentRestartTest.java
@@ -18,46 +18,59 @@
package org.apache.flink.kubernetes.operator.controller;
import org.apache.flink.configuration.Configuration;
+import org.apache.flink.kubernetes.operator.OperatorTestBase;
import org.apache.flink.kubernetes.operator.TestUtils;
-import org.apache.flink.kubernetes.operator.TestingFlinkService;
import org.apache.flink.kubernetes.operator.api.FlinkDeployment;
import org.apache.flink.kubernetes.operator.api.spec.FlinkVersion;
import org.apache.flink.kubernetes.operator.api.spec.UpgradeMode;
+import org.apache.flink.kubernetes.operator.api.status.CheckpointInfo;
+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.SnapshotTriggerType;
import org.apache.flink.kubernetes.operator.config.FlinkConfigManager;
+import org.apache.flink.kubernetes.operator.observer.SnapshotObserver;
+import org.apache.flink.kubernetes.operator.utils.SnapshotStatus;
+import org.apache.flink.kubernetes.operator.utils.SnapshotUtils;
+import org.apache.flink.runtime.jobgraph.SavepointConfigOptions;
import io.fabric8.kubernetes.client.KubernetesClient;
import io.fabric8.kubernetes.client.server.mock.EnableKubernetesMockClient;
import io.javaoperatorsdk.operator.api.reconciler.Context;
+import lombok.Getter;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.EnumSource;
import org.junit.jupiter.params.provider.MethodSource;
import static
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions.OPERATOR_JOB_RESTART_FAILED;
+import static
org.apache.flink.kubernetes.operator.reconciler.SnapshotType.CHECKPOINT;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
/**
* @link Unhealthy deployment restart tests
*/
@EnableKubernetesMockClient(crud = true)
-public class FailedDeploymentRestartTest {
+public class FailedDeploymentRestartTest extends OperatorTestBase {
private FlinkConfigManager configManager;
- private TestingFlinkService flinkService;
private Context<FlinkDeployment> context;
private TestingFlinkDeploymentController testController;
+ private SnapshotObserver<FlinkDeployment, FlinkDeploymentStatus> observer;
- private KubernetesClient kubernetesClient;
+ @Getter private KubernetesClient kubernetesClient;
@BeforeEach
public void setup() {
var configuration = new Configuration();
configuration.set(OPERATOR_JOB_RESTART_FAILED, true);
configManager = new FlinkConfigManager(configuration);
- flinkService = new TestingFlinkService(kubernetesClient);
context = flinkService.getContext();
testController = new TestingFlinkDeploymentController(configManager,
flinkService);
kubernetesClient.resource(TestUtils.buildApplicationCluster()).createOrReplace();
+ observer = new SnapshotObserver<>(eventRecorder);
}
@ParameterizedTest
@@ -98,4 +111,68 @@ public class FailedDeploymentRestartTest {
appCluster.getSpec(),
appCluster.getStatus().getReconciliationStatus().deserializeLastReconciledSpec());
}
+
+ @ParameterizedTest
+ @EnumSource(UpgradeMode.class)
+ public void verifyFailedApplicationRecoveryWithCheckpoint(UpgradeMode
upgradeMode)
+ throws Exception {
+ FlinkDeployment appCluster = TestUtils.buildApplicationCluster();
+ appCluster.getSpec().getJob().setUpgradeMode(upgradeMode);
+
+ // Start a healthy deployment
+ testController.reconcile(appCluster, context);
+ testController.reconcile(appCluster, context);
+ testController.reconcile(appCluster, context);
+
+ // Mark job_id
+ String jobId = appCluster.getStatus().getJobStatus().getJobId();
+ assertNotNull(jobId);
+ assertEquals(
+ JobManagerDeploymentStatus.READY,
+ appCluster.getStatus().getJobManagerDeploymentStatus());
+ assertEquals("RUNNING",
appCluster.getStatus().getJobStatus().getState());
+
assertNull(flinkService.getSubmittedConf().get(SavepointConfigOptions.SAVEPOINT_PATH));
+
+ // trigger checkpoint
+ CheckpointInfo checkpointInfo =
appCluster.getStatus().getJobStatus().getCheckpointInfo();
+ flinkService.triggerCheckpoint(
+ null,
+ SnapshotTriggerType.PERIODIC,
+ checkpointInfo,
+ configManager.getObserveConfig(appCluster));
+
+ // Pending
+ observer.observeCheckpointStatus(getResourceContext(appCluster));
+ // Completed
+ observer.observeCheckpointStatus(getResourceContext(appCluster));
+
assertFalse(SnapshotUtils.checkpointInProgress(appCluster.getStatus().getJobStatus()));
+ assertEquals(
+ SnapshotUtils.getLastSnapshotStatus(appCluster, CHECKPOINT),
+ SnapshotStatus.SUCCEEDED);
+
+ // Make deployment unhealthy
+ flinkService.markApplicationJobFailedWithError(
+ flinkService.listJobs().get(0).f1.getJobId(), "Failed job");
+ testController.reconcile(appCluster, context);
+ assertEquals(
+ JobManagerDeploymentStatus.DEPLOYING,
+ appCluster.getStatus().getJobManagerDeploymentStatus());
+
+ // After restart the deployment is healthy again
+ testController.reconcile(appCluster, context);
+ testController.reconcile(appCluster, context);
+ assertEquals(
+ JobManagerDeploymentStatus.READY,
+ appCluster.getStatus().getJobManagerDeploymentStatus());
+ assertEquals("RUNNING",
appCluster.getStatus().getJobStatus().getState());
+
+ // check savepoint_path
+ if (upgradeMode != UpgradeMode.STATELESS) {
+ assertEquals(
+
flinkService.getSubmittedConf().get(SavepointConfigOptions.SAVEPOINT_PATH),
+ "ck_0");
+ } else {
+
assertNull(flinkService.getSubmittedConf().get(SavepointConfigOptions.SAVEPOINT_PATH));
+ }
+ }
}