From 21867b8e8e4adefd449e93dde9c0a0857f2b25ec Mon Sep 17 00:00:00 2001 From: Yanis Djeridi Date: Mon, 30 Mar 2026 10:37:16 -0400 Subject: [PATCH 1/4] 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 2/4] 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 3/4] 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 4/4] 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); }