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..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 @@ -89,6 +89,12 @@ public UpdateControl reconcile( @Override public DeleteControl cleanup( FlinkStateSnapshot flinkStateSnapshot, Context josdkContext) { + 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); 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..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 @@ -704,6 +704,20 @@ 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); + } + private FlinkStateSnapshot createSavepoint(FlinkDeployment deployment) { return createSavepoint(deployment, false, 7); }