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