From 216a86ec74b9bcc5f53b3ba270422a899e5f0a37 Mon Sep 17 00:00:00 2001 From: Jennifer Xiong Date: Wed, 18 Mar 2026 23:32:47 -0700 Subject: [PATCH 01/11] fix stale trigger point id bug --- .../controller/bluegreen/BlueGreenDeploymentService.java | 1 + .../controller/FlinkBlueGreenDeploymentControllerTest.java | 3 +++ 2 files changed, 4 insertions(+) diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java index 9543fecabd..931bf36e5c 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java @@ -685,6 +685,7 @@ private UpdateControl abortDeployment( FlinkBlueGreenDeploymentState previousState = getPreviousState(nextState, context.getDeployments()); context.getDeploymentStatus().setBlueGreenState(previousState); + context.getDeploymentStatus().setSavepointTriggerId(null); var error = String.format( diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java index 528b61727d..1aa021ecec 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java @@ -546,6 +546,9 @@ public void verifyFailureDuringTransition(FlinkVersion flinkVersion) throws Exce ReconciliationState.UPGRADING, flinkDeployments.get(1).getStatus().getReconciliationStatus().getState()); assertTrue(instantStrToMillis(rs.reconciledStatus.getAbortTimestamp()) > 0); + // savepointTriggerId must be cleared on abort so the next transition + // triggers a fresh savepoint instead of reusing a stale triggerId + assertNull(rs.reconciledStatus.getSavepointTriggerId()); // Simulate another change in the spec to trigger a redeployment customValue = UUID.randomUUID().toString(); From fb4d088b8623adbada5da18f58b767fdcd94fb9a Mon Sep 17 00:00:00 2001 From: Jennifer Xiong Date: Thu, 19 Mar 2026 00:04:33 -0700 Subject: [PATCH 02/11] last reconciled bug --- .../bluegreen/BlueGreenDeploymentService.java | 22 ++++++++++++++----- ...linkBlueGreenDeploymentControllerTest.java | 7 ++++++ 2 files changed, 23 insertions(+), 6 deletions(-) diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java index 9543fecabd..1441ba9797 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java @@ -162,10 +162,18 @@ public UpdateControl checkAndInitiateDeployment( .rescheduleAfter(getReconciliationReschedInterval(context)); } - setLastReconciledSpec(context); try { - return startTransition( - context, currentBlueGreenDeploymentType, currentFlinkDeployment); + var result = + startTransition( + context, + currentBlueGreenDeploymentType, + currentFlinkDeployment); + // Only stamp lastReconciledSpec after the transition + // succeeds. If stamped before and the transition + // fails/aborts, lastReconciledSpec drifts from the + // active child's actual spec + setLastReconciledSpec(context); + return result; } catch (Exception e) { var error = "Could not start Transition. Details: " + e.getMessage(); context.getDeploymentStatus().setSavepointTriggerId(null); @@ -180,10 +188,12 @@ public UpdateControl checkAndInitiateDeployment( "Savepoint redeploy triggered for '{}', using initialSavepointPath: {}", context.getBgDeployment().getMetadata().getName(), Objects.toString(jobSpec.getInitialSavepointPath(), "")); - setLastReconciledSpec(context); try { - return startSavepointRedeployTransition( - context, currentBlueGreenDeploymentType); + var result = + startSavepointRedeployTransition( + context, currentBlueGreenDeploymentType); + setLastReconciledSpec(context); + return result; } catch (Exception e) { var error = "Could not start Savepoint Redeploy Transition. Details: " diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java index 528b61727d..8b247c3dbe 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java @@ -767,6 +767,13 @@ public void verifySavepointFetchFailureRecovery(FlinkVersion flinkVersion) throw rs = reconcile(rs.deployment); assertFailingWithError(rs, "Could not start Transition", error); + // lastReconciledSpec must NOT contain the failed spec change — otherwise + // subsequent deploys see no diff and fall into PATCH_CHILD (in-place + // upgrade) instead of a proper B/G transition + assertFalse( + rs.reconciledStatus.getLastReconciledSpec().contains(customValue), + "lastReconciledSpec should not be stamped when transition fails"); + // Recovery: Clear the fetch error and try again with new spec change flinkService.clearSavepointFetchError(); customValue = UUID.randomUUID().toString() + "_recovery"; From 7b2a294d7f6ad6b21c8adc02475fd9495508d531 Mon Sep 17 00:00:00 2001 From: Yanis Djeridi Date: Fri, 20 Mar 2026 00:12:49 -0400 Subject: [PATCH 03/11] Revert "Move setLastReconciledSpec to after successful transition start" --- .../bluegreen/BlueGreenDeploymentService.java | 22 +++++-------------- ...linkBlueGreenDeploymentControllerTest.java | 7 ------ 2 files changed, 6 insertions(+), 23 deletions(-) diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java index 1441ba9797..9543fecabd 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java @@ -162,18 +162,10 @@ public UpdateControl checkAndInitiateDeployment( .rescheduleAfter(getReconciliationReschedInterval(context)); } + setLastReconciledSpec(context); try { - var result = - startTransition( - context, - currentBlueGreenDeploymentType, - currentFlinkDeployment); - // Only stamp lastReconciledSpec after the transition - // succeeds. If stamped before and the transition - // fails/aborts, lastReconciledSpec drifts from the - // active child's actual spec - setLastReconciledSpec(context); - return result; + return startTransition( + context, currentBlueGreenDeploymentType, currentFlinkDeployment); } catch (Exception e) { var error = "Could not start Transition. Details: " + e.getMessage(); context.getDeploymentStatus().setSavepointTriggerId(null); @@ -188,12 +180,10 @@ public UpdateControl checkAndInitiateDeployment( "Savepoint redeploy triggered for '{}', using initialSavepointPath: {}", context.getBgDeployment().getMetadata().getName(), Objects.toString(jobSpec.getInitialSavepointPath(), "")); + setLastReconciledSpec(context); try { - var result = - startSavepointRedeployTransition( - context, currentBlueGreenDeploymentType); - setLastReconciledSpec(context); - return result; + return startSavepointRedeployTransition( + context, currentBlueGreenDeploymentType); } catch (Exception e) { var error = "Could not start Savepoint Redeploy Transition. Details: " diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java index 8b247c3dbe..528b61727d 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java @@ -767,13 +767,6 @@ public void verifySavepointFetchFailureRecovery(FlinkVersion flinkVersion) throw rs = reconcile(rs.deployment); assertFailingWithError(rs, "Could not start Transition", error); - // lastReconciledSpec must NOT contain the failed spec change — otherwise - // subsequent deploys see no diff and fall into PATCH_CHILD (in-place - // upgrade) instead of a proper B/G transition - assertFalse( - rs.reconciledStatus.getLastReconciledSpec().contains(customValue), - "lastReconciledSpec should not be stamped when transition fails"); - // Recovery: Clear the fetch error and try again with new spec change flinkService.clearSavepointFetchError(); customValue = UUID.randomUUID().toString() + "_recovery"; From b31e482dc7a17269256ea3cb7f493a60e793c88d Mon Sep 17 00:00:00 2001 From: Yanis Djeridi Date: Fri, 20 Mar 2026 13:57:28 -0400 Subject: [PATCH 04/11] last reconciled bug fix --- .../bluegreen/BlueGreenDeploymentService.java | 22 ++++++++++++++----- ...linkBlueGreenDeploymentControllerTest.java | 7 ++++++ 2 files changed, 23 insertions(+), 6 deletions(-) diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java index 931bf36e5c..6274881df2 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java @@ -162,10 +162,18 @@ public UpdateControl checkAndInitiateDeployment( .rescheduleAfter(getReconciliationReschedInterval(context)); } - setLastReconciledSpec(context); try { - return startTransition( - context, currentBlueGreenDeploymentType, currentFlinkDeployment); + var result = + startTransition( + context, + currentBlueGreenDeploymentType, + currentFlinkDeployment); + // Only stamp lastReconciledSpec after the transition + // succeeds. If stamped before and the transition + // fails/aborts, lastReconciledSpec drifts from the + // active child's actual spec + setLastReconciledSpec(context); + return result; } catch (Exception e) { var error = "Could not start Transition. Details: " + e.getMessage(); context.getDeploymentStatus().setSavepointTriggerId(null); @@ -180,10 +188,12 @@ public UpdateControl checkAndInitiateDeployment( "Savepoint redeploy triggered for '{}', using initialSavepointPath: {}", context.getBgDeployment().getMetadata().getName(), Objects.toString(jobSpec.getInitialSavepointPath(), "")); - setLastReconciledSpec(context); try { - return startSavepointRedeployTransition( - context, currentBlueGreenDeploymentType); + var result = + startSavepointRedeployTransition( + context, currentBlueGreenDeploymentType); + setLastReconciledSpec(context); + return result; } catch (Exception e) { var error = "Could not start Savepoint Redeploy Transition. Details: " diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java index 1aa021ecec..e657bb7227 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java @@ -770,6 +770,13 @@ public void verifySavepointFetchFailureRecovery(FlinkVersion flinkVersion) throw rs = reconcile(rs.deployment); assertFailingWithError(rs, "Could not start Transition", error); + // lastReconciledSpec must NOT contain the failed spec change — otherwise + // subsequent deploys see no diff and fall into PATCH_CHILD (in-place + // upgrade) instead of a proper B/G transition + assertFalse( + rs.reconciledStatus.getLastReconciledSpec().contains(customValue), + "lastReconciledSpec should not be stamped when transition fails"); + // Recovery: Clear the fetch error and try again with new spec change flinkService.clearSavepointFetchError(); customValue = UUID.randomUUID().toString() + "_recovery"; From 7c617c22fa15cf4c9dbb5aae1c91e10b245f52f8 Mon Sep 17 00:00:00 2001 From: Yanis Djeridi Date: Fri, 20 Mar 2026 14:02:38 -0400 Subject: [PATCH 05/11] match changes --- .../controller/FlinkBlueGreenDeploymentControllerTest.java | 3 --- 1 file changed, 3 deletions(-) diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java index e657bb7227..8b247c3dbe 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java @@ -546,9 +546,6 @@ public void verifyFailureDuringTransition(FlinkVersion flinkVersion) throws Exce ReconciliationState.UPGRADING, flinkDeployments.get(1).getStatus().getReconciliationStatus().getState()); assertTrue(instantStrToMillis(rs.reconciledStatus.getAbortTimestamp()) > 0); - // savepointTriggerId must be cleared on abort so the next transition - // triggers a fresh savepoint instead of reusing a stale triggerId - assertNull(rs.reconciledStatus.getSavepointTriggerId()); // Simulate another change in the spec to trigger a redeployment customValue = UUID.randomUUID().toString(); From 80fe58b004577b2502d67d28238ac62389ea0ba3 Mon Sep 17 00:00:00 2001 From: Yanis Djeridi Date: Fri, 20 Mar 2026 14:04:15 -0400 Subject: [PATCH 06/11] match --- .../controller/bluegreen/BlueGreenDeploymentService.java | 1 - 1 file changed, 1 deletion(-) diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java index 6274881df2..1441ba9797 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java @@ -695,7 +695,6 @@ private UpdateControl abortDeployment( FlinkBlueGreenDeploymentState previousState = getPreviousState(nextState, context.getDeployments()); context.getDeploymentStatus().setBlueGreenState(previousState); - context.getDeploymentStatus().setSavepointTriggerId(null); var error = String.format( From fe9cc1f05df3cda0dbb250af0b7d0de90d47d9d6 Mon Sep 17 00:00:00 2001 From: Yanis Djeridi Date: Fri, 20 Mar 2026 14:16:05 -0400 Subject: [PATCH 07/11] Restore PR #23 changes that were accidentally removed Made-with: Cursor --- .../controller/bluegreen/BlueGreenDeploymentService.java | 1 + .../controller/FlinkBlueGreenDeploymentControllerTest.java | 3 +++ 2 files changed, 4 insertions(+) diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java index 1441ba9797..6274881df2 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java @@ -695,6 +695,7 @@ private UpdateControl abortDeployment( FlinkBlueGreenDeploymentState previousState = getPreviousState(nextState, context.getDeployments()); context.getDeploymentStatus().setBlueGreenState(previousState); + context.getDeploymentStatus().setSavepointTriggerId(null); var error = String.format( diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java index 8b247c3dbe..e657bb7227 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java @@ -546,6 +546,9 @@ public void verifyFailureDuringTransition(FlinkVersion flinkVersion) throws Exce ReconciliationState.UPGRADING, flinkDeployments.get(1).getStatus().getReconciliationStatus().getState()); assertTrue(instantStrToMillis(rs.reconciledStatus.getAbortTimestamp()) > 0); + // savepointTriggerId must be cleared on abort so the next transition + // triggers a fresh savepoint instead of reusing a stale triggerId + assertNull(rs.reconciledStatus.getSavepointTriggerId()); // Simulate another change in the spec to trigger a redeployment customValue = UUID.randomUUID().toString(); From 21867b8e8e4adefd449e93dde9c0a0857f2b25ec Mon Sep 17 00:00:00 2001 From: Yanis Djeridi Date: Mon, 30 Mar 2026 10:37:16 -0400 Subject: [PATCH 08/11] orphaned state snapshot fix --- .../operator/controller/FlinkStateSnapshotController.java | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotController.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotController.java index 8c36dc0efd..ee0b2ec2ee 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotController.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotController.java @@ -89,6 +89,9 @@ public UpdateControl reconcile( @Override public DeleteControl cleanup( FlinkStateSnapshot flinkStateSnapshot, Context josdkContext) { + flinkStateSnapshot.setStatus( + Objects.requireNonNullElseGet( + flinkStateSnapshot.getStatus(), FlinkStateSnapshotStatus::new)); var ctx = ctxFactory.getFlinkStateSnapshotContext(flinkStateSnapshot, josdkContext); try { metricManager.onRemove(flinkStateSnapshot); @@ -113,6 +116,8 @@ public DeleteControl cleanup( @Override public ErrorStatusUpdateControl updateErrorStatus( FlinkStateSnapshot resource, Context context, Exception e) { + resource.setStatus( + Objects.requireNonNullElseGet(resource.getStatus(), FlinkStateSnapshotStatus::new)); var ctx = ctxFactory.getFlinkStateSnapshotContext(resource, context); ReconciliationUtils.updateForReconciliationError(ctx, e); From c7bd5622a0f2f3409323644a841df3f261242a0c Mon Sep 17 00:00:00 2001 From: Yanis Djeridi Date: Mon, 30 Mar 2026 14:31:13 -0400 Subject: [PATCH 09/11] add tests --- .../FlinkStateSnapshotControllerTest.java | 27 +++++++++++++++++++ 1 file changed, 27 insertions(+) diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotControllerTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotControllerTest.java index cfe6947343..e0647957d9 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotControllerTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotControllerTest.java @@ -704,6 +704,33 @@ public void testMetrics() { assertSnapshotMetrics(listener, TestUtils.TEST_NAMESPACE, Map.of(), Map.of()); } + @Test + public void testCleanupWithNullStatus() { + var deployment = createDeployment(); + context = TestUtils.createSnapshotContext(client, deployment); + + var savepoint = createSavepoint(deployment); + savepoint.setStatus(null); + assertDeleteControl(controller.cleanup(savepoint, context), true, null); + + var checkpoint = createCheckpoint(deployment, CheckpointType.FULL, 0); + checkpoint.setStatus(null); + assertDeleteControl(controller.cleanup(checkpoint, context), true, null); + } + + @Test + public void testUpdateErrorStatusWithNullStatus() { + var deployment = createDeployment(); + context = TestUtils.createSnapshotContext(client, deployment); + var snapshot = createSavepoint(deployment); + snapshot.setStatus(null); + + controller.updateErrorStatus(snapshot, context, new Exception("test error")); + + assertThat(snapshot.getStatus()).isNotNull(); + assertThat(snapshot.getStatus().getError()).isNotNull(); + } + private FlinkStateSnapshot createSavepoint(FlinkDeployment deployment) { return createSavepoint(deployment, false, 7); } From 49a9610e22f676fac60a165b2ff5c8240184945f Mon Sep 17 00:00:00 2001 From: Yanis Djeridi Date: Mon, 30 Mar 2026 15:08:52 -0400 Subject: [PATCH 10/11] change approach --- .../controller/FlinkStateSnapshotController.java | 14 +++++++++----- 1 file changed, 9 insertions(+), 5 deletions(-) diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotController.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotController.java index ee0b2ec2ee..e43d77f7be 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotController.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotController.java @@ -89,9 +89,12 @@ public UpdateControl reconcile( @Override public DeleteControl cleanup( FlinkStateSnapshot flinkStateSnapshot, Context josdkContext) { - flinkStateSnapshot.setStatus( - Objects.requireNonNullElseGet( - flinkStateSnapshot.getStatus(), FlinkStateSnapshotStatus::new)); + if (flinkStateSnapshot.getStatus() == null) { + LOG.info( + "Snapshot {} has no status, was never reconciled. Removing finalizer.", + flinkStateSnapshot.getMetadata().getName()); + return DeleteControl.defaultDelete(); + } var ctx = ctxFactory.getFlinkStateSnapshotContext(flinkStateSnapshot, josdkContext); try { metricManager.onRemove(flinkStateSnapshot); @@ -116,8 +119,9 @@ public DeleteControl cleanup( @Override public ErrorStatusUpdateControl updateErrorStatus( FlinkStateSnapshot resource, Context context, Exception e) { - resource.setStatus( - Objects.requireNonNullElseGet(resource.getStatus(), FlinkStateSnapshotStatus::new)); + if (resource.getStatus() == null) { + resource.setStatus(new FlinkStateSnapshotStatus()); + } var ctx = ctxFactory.getFlinkStateSnapshotContext(resource, context); ReconciliationUtils.updateForReconciliationError(ctx, e); From 7f7e2a6829ddce5fb3e419702dabfdca6895665f Mon Sep 17 00:00:00 2001 From: Yanis Djeridi Date: Mon, 30 Mar 2026 16:20:40 -0400 Subject: [PATCH 11/11] Remove updateErrorStatus check --- .../controller/FlinkStateSnapshotController.java | 3 --- .../FlinkStateSnapshotControllerTest.java | 13 ------------- 2 files changed, 16 deletions(-) diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotController.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotController.java index e43d77f7be..57626c7865 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotController.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotController.java @@ -119,9 +119,6 @@ public DeleteControl cleanup( @Override public ErrorStatusUpdateControl updateErrorStatus( FlinkStateSnapshot resource, Context context, Exception e) { - if (resource.getStatus() == null) { - resource.setStatus(new FlinkStateSnapshotStatus()); - } var ctx = ctxFactory.getFlinkStateSnapshotContext(resource, context); ReconciliationUtils.updateForReconciliationError(ctx, e); diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotControllerTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotControllerTest.java index e0647957d9..68b979625a 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotControllerTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkStateSnapshotControllerTest.java @@ -718,19 +718,6 @@ public void testCleanupWithNullStatus() { assertDeleteControl(controller.cleanup(checkpoint, context), true, null); } - @Test - public void testUpdateErrorStatusWithNullStatus() { - var deployment = createDeployment(); - context = TestUtils.createSnapshotContext(client, deployment); - var snapshot = createSavepoint(deployment); - snapshot.setStatus(null); - - controller.updateErrorStatus(snapshot, context, new Exception("test error")); - - assertThat(snapshot.getStatus()).isNotNull(); - assertThat(snapshot.getStatus().getError()).isNotNull(); - } - private FlinkStateSnapshot createSavepoint(FlinkDeployment deployment) { return createSavepoint(deployment, false, 7); }