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 ac2245a4 [FLINK-39356] Fix NPE during FlinkStateSnapshot cleanup when
status is null (#1169)
ac2245a4 is described below
commit ac2245a4dc34bfce83c2ca7a59616cd42b738b1f
Author: Nihar Rao <[email protected]>
AuthorDate: Wed Aug 5 23:44:05 2026 -0400
[FLINK-39356] Fix NPE during FlinkStateSnapshot cleanup when status is null
(#1169)
---
.../operator/controller/FlinkStateSnapshotController.java | 3 +++
.../controller/FlinkStateSnapshotControllerTest.java | 14 ++++++++++++++
2 files changed, 17 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 1bb2a3d0..990acdfd 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 class FlinkStateSnapshotController
@Override
public DeleteControl cleanup(
FlinkStateSnapshot flinkStateSnapshot, Context<FlinkStateSnapshot>
josdkContext) {
+ if (flinkStateSnapshot.getStatus() == null) {
+ 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 c23478ad..433a2f82 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
@@ -729,6 +729,20 @@ public class FlinkStateSnapshotControllerTest {
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);
}