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));
+        }
+    }
 }

Reply via email to