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 0ee51116 [FLINK-40162] Handle UpgradeFailureException in 
FlinkSessionJobController (#1166)
0ee51116 is described below

commit 0ee511164d604930f0170a191b4b24a613b4586b
Author: Milind L <[email protected]>
AuthorDate: Mon Aug 3 11:52:01 2026 +0530

    [FLINK-40162] Handle UpgradeFailureException in FlinkSessionJobController 
(#1166)
    
    Move observer.observe() inside try-catch and add explicit catch for
    UpgradeFailureException, mirroring FlinkDeploymentController's pattern.
    This prevents FlinkSessionJob from getting stuck in RECONCILING state
    when a job finishes quickly with checkpointing enabled but no
    state.checkpoints.dir configured.
---
 .../controller/FlinkSessionJobController.java      | 23 ++++----
 .../controller/FlinkSessionJobControllerTest.java  | 64 ++++++++++++++++++++++
 2 files changed, 77 insertions(+), 10 deletions(-)

diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/FlinkSessionJobController.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/FlinkSessionJobController.java
index d5f0f867..3fc3df62 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/FlinkSessionJobController.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/FlinkSessionJobController.java
@@ -24,6 +24,7 @@ import 
org.apache.flink.kubernetes.operator.api.status.FlinkSessionJobStatus;
 import 
org.apache.flink.kubernetes.operator.api.validation.FlinkResourceValidator;
 import org.apache.flink.kubernetes.operator.config.FlinkConfigManager;
 import org.apache.flink.kubernetes.operator.exception.ReconciliationException;
+import org.apache.flink.kubernetes.operator.exception.UpgradeFailureException;
 import org.apache.flink.kubernetes.operator.health.CanaryResourceManager;
 import org.apache.flink.kubernetes.operator.observer.Observer;
 import org.apache.flink.kubernetes.operator.reconciler.Reconciler;
@@ -107,18 +108,20 @@ public class FlinkSessionJobController
             return UpdateControl.noUpdate();
         }
 
-        observer.observe(ctx);
-        if (!validateSessionJob(ctx)) {
-            statusRecorder.patchAndCacheStatus(flinkSessionJob, 
ctx.getKubernetesClient());
-            return ReconciliationUtils.toUpdateControl(
-                    ctx.getOperatorConfig(), flinkSessionJob, previousJob, 
false);
-        }
-
         try {
+            observer.observe(ctx);
+            if (!validateSessionJob(ctx)) {
+                statusRecorder.patchAndCacheStatus(flinkSessionJob, 
ctx.getKubernetesClient());
+                return ReconciliationUtils.toUpdateControl(
+                        ctx.getOperatorConfig(), flinkSessionJob, previousJob, 
false);
+            }
             statusRecorder.patchAndCacheStatus(flinkSessionJob, 
ctx.getKubernetesClient());
             reconciler.reconcile(ctx);
+        } catch (UpgradeFailureException ufe) {
+            ReconciliationUtils.updateForReconciliationError(ctx, ufe);
+            triggerErrorEvent(ctx, ufe, ufe.getReason());
         } catch (Exception e) {
-            triggerErrorEvent(ctx, e);
+            triggerErrorEvent(ctx, e, EventRecorder.Reason.Error.name());
             throw new ReconciliationException(e);
         }
         statusRecorder.patchAndCacheStatus(flinkSessionJob, 
ctx.getKubernetesClient());
@@ -158,11 +161,11 @@ public class FlinkSessionJobController
         return deleteControl;
     }
 
-    private void triggerErrorEvent(FlinkResourceContext<?> ctx, Exception e) {
+    private void triggerErrorEvent(FlinkResourceContext<?> ctx, Exception e, 
String reason) {
         eventRecorder.triggerEvent(
                 ctx.getResource(),
                 EventRecorder.Type.Warning,
-                EventRecorder.Reason.Error.name(),
+                reason,
                 ExceptionUtils.getExceptionMessage(e),
                 EventRecorder.Component.Job,
                 ctx.getKubernetesClient());
diff --git 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkSessionJobControllerTest.java
 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkSessionJobControllerTest.java
index 49edb5df..af4302be 100644
--- 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkSessionJobControllerTest.java
+++ 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkSessionJobControllerTest.java
@@ -36,6 +36,7 @@ import 
org.apache.flink.kubernetes.operator.observer.JobStatusObserver;
 import org.apache.flink.kubernetes.operator.service.CheckpointHistoryWrapper;
 import org.apache.flink.kubernetes.operator.utils.EventRecorder;
 import org.apache.flink.runtime.client.JobStatusMessage;
+import 
org.apache.flink.runtime.state.memory.NonPersistentMetadataCheckpointStorageLocation;
 import org.apache.flink.util.SerializedThrowable;
 
 import io.fabric8.kubernetes.client.KubernetesClient;
@@ -54,6 +55,7 @@ import java.util.Optional;
 import java.util.stream.Collectors;
 
 import static org.apache.flink.api.common.JobStatus.CANCELLING;
+import static org.apache.flink.api.common.JobStatus.FINISHED;
 import static org.apache.flink.api.common.JobStatus.RECONCILING;
 import static org.apache.flink.api.common.JobStatus.RUNNING;
 import static 
org.apache.flink.kubernetes.operator.TestUtils.MAX_RECONCILE_TIMES;
@@ -831,4 +833,66 @@ class FlinkSessionJobControllerTest {
                 
sessionJob.getStatus().getReconciliationStatus().getLastReconciledSpec(),
                 
sessionJob.getStatus().getReconciliationStatus().getLastStableSpec());
     }
+
+    @Test
+    public void testJobFinishedWithNonExternallyAddressableCheckpoint() throws 
Exception {
+        // Submit the session job
+        testController.reconcile(sessionJob, context);
+        assertEquals(RECONCILING, 
sessionJob.getStatus().getJobStatus().getState());
+
+        // Observe RUNNING state
+        testController.reconcile(sessionJob, context);
+        assertEquals(RUNNING, 
sessionJob.getStatus().getJobStatus().getState());
+
+        // Simulate the job finishing quickly
+        var jobs = flinkService.listJobs();
+        assertEquals(1, jobs.size());
+        var job = jobs.get(0);
+        jobs.set(
+                0,
+                Tuple3.of(
+                        job.f0,
+                        new JobStatusMessage(
+                                job.f1.getJobId(),
+                                job.f1.getJobName(),
+                                FINISHED,
+                                job.f1.getStartTime()),
+                        job.f2));
+
+        // Simulate checkpointing enabled but no state.checkpoints.dir 
configured
+        flinkService.setCheckpointInfo(
+                Tuple2.of(
+                        Optional.of(
+                                new 
CheckpointHistoryWrapper.CompletedCheckpointInfo(
+                                        1L,
+                                        
NonPersistentMetadataCheckpointStorageLocation
+                                                .EXTERNAL_POINTER,
+                                        System.currentTimeMillis())),
+                        Optional.empty()));
+
+        // Reconcile - should handle the UpgradeFailureException gracefully
+        testController.events().clear();
+        testController.reconcile(sessionJob, context);
+
+        // Job status should reflect FINISHED, not be stuck at RECONCILING
+        assertEquals(FINISHED, 
sessionJob.getStatus().getJobStatus().getState());
+
+        // Reconciliation status should reflect the error
+        assertEquals(
+                ReconciliationState.DEPLOYED,
+                sessionJob.getStatus().getReconciliationStatus().getState());
+
+        // Error should be recorded about the non-externally-addressable 
checkpoint
+        assertNotNull(sessionJob.getStatus().getError());
+        assertTrue(sessionJob.getStatus().getError().contains("not externally 
addressable"));
+
+        // Verify warning event was emitted with the specific reason
+        var errorEvent =
+                testController.events().stream()
+                        .filter(e -> 
e.getType().equals(EventRecorder.Type.Warning.toString()))
+                        .findFirst();
+        assertTrue(errorEvent.isPresent());
+        assertEquals("CheckpointNotFound", errorEvent.get().getReason());
+        assertTrue(errorEvent.get().getMessage().contains("not externally 
addressable"));
+    }
 }

Reply via email to