This is an automated email from the ASF dual-hosted git repository.
mbalassi 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 e73363f3 [FLINK-35278] Fix occasional NPE on getting latest resource
for status replace
e73363f3 is described below
commit e73363f3486ed9e1df5cc05c9d0baec7c8c3a37f
Author: Ferenc Csaky <[email protected]>
AuthorDate: Sat May 4 11:12:40 2024 +0200
[FLINK-35278] Fix occasional NPE on getting latest resource for status
replace
---
.../kubernetes/operator/utils/StatusRecorder.java | 87 +++++++++++++---------
.../operator/utils/StatusRecorderTest.java | 27 ++++++-
2 files changed, 76 insertions(+), 38 deletions(-)
diff --git
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/StatusRecorder.java
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/StatusRecorder.java
index 0197dadd..1b7e3fc2 100644
---
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/StatusRecorder.java
+++
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/StatusRecorder.java
@@ -18,6 +18,7 @@
package org.apache.flink.kubernetes.operator.utils;
+import org.apache.flink.annotation.VisibleForTesting;
import org.apache.flink.kubernetes.operator.api.AbstractFlinkResource;
import org.apache.flink.kubernetes.operator.api.FlinkDeployment;
import
org.apache.flink.kubernetes.operator.api.lifecycle.ResourceLifecycleState;
@@ -127,47 +128,65 @@ public class StatusRecorder<
} catch (KubernetesClientException kce) {
// 409 is the error code for conflicts resulting from the
locking
if (kce.getCode() == 409) {
- var currentVersion =
resource.getMetadata().getResourceVersion();
- LOG.debug(
- "Could not apply status update for resource
version {}",
- currentVersion);
-
- var latest = client.resource(resource).get();
- var latestVersion =
latest.getMetadata().getResourceVersion();
-
- if (latestVersion.equals(currentVersion)) {
- // This should not happen as long as the client works
consistently
- LOG.error("Unable to fetch latest resource version");
- throw kce;
- }
-
- if (latest.getStatus().equals(prevStatus)) {
- if (retries++ < 3) {
- LOG.debug(
- "Retrying status update for latest version
{}", latestVersion);
-
resource.getMetadata().setResourceVersion(latestVersion);
- } else {
- // If we cannot get the latest version in 3 tries
we throw the error to
- // retry with delay
- throw kce;
- }
- } else {
- throw new StatusConflictException(
- "Status have been modified externally in
version "
- + latestVersion
- + " Previous: "
- +
objectMapper.writeValueAsString(prevStatus)
- + " Latest: "
- +
objectMapper.writeValueAsString(latest.getStatus()));
- }
+ handleLockingError(resource, prevStatus, client, retries,
kce);
+ ++retries;
} else {
- // We simply throw non conflict errors, to trigger retry
with delay
+ // We simply throw non-conflict errors, to trigger retry
with delay
throw kce;
}
}
}
}
+ @VisibleForTesting
+ void handleLockingError(
+ CR resource,
+ STATUS prevStatus,
+ KubernetesClient client,
+ int retries,
+ KubernetesClientException kce)
+ throws JsonProcessingException {
+
+ var currentVersion = resource.getMetadata().getResourceVersion();
+ LOG.debug("Could not apply status update for resource version {}",
currentVersion);
+
+ var latest = client.resource(resource).get();
+ if (latest == null || latest.getMetadata() == null) {
+ // This can happen occasionally, we throw the error to retry with
delay.
+ throw new KubernetesClientException(
+ String.format(
+ "Failed to retrieve latest %s",
+ latest == null ? "resource" : "metadata"),
+ kce);
+ }
+
+ var latestVersion = latest.getMetadata().getResourceVersion();
+ if (currentVersion.equals(latestVersion)) {
+ // This should not happen as long as the client works consistently
+ LOG.error("Unable to fetch latest resource version");
+ throw kce;
+ }
+
+ if (latest.getStatus().equals(prevStatus)) {
+ if (retries < 3) {
+ LOG.debug("Retrying status update for latest version {}",
latestVersion);
+ resource.getMetadata().setResourceVersion(latestVersion);
+ } else {
+ // If we cannot get the latest version in 3 tries we throw the
error to
+ // retry with delay
+ throw kce;
+ }
+ } else {
+ throw new StatusConflictException(
+ "Status have been modified externally in version "
+ + latestVersion
+ + " Previous: "
+ + objectMapper.writeValueAsString(prevStatus)
+ + " Latest: "
+ +
objectMapper.writeValueAsString(latest.getStatus()));
+ }
+ }
+
/**
* Update the custom resource status based on the in-memory cached to
ensure that any status
* updates that we made previously are always visible in the
reconciliation loop. This is
diff --git
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/StatusRecorderTest.java
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/StatusRecorderTest.java
index aab8eb80..6f9b2e52 100644
---
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/StatusRecorderTest.java
+++
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/StatusRecorderTest.java
@@ -24,11 +24,13 @@ import
org.apache.flink.kubernetes.operator.api.status.ReconciliationState;
import org.apache.flink.kubernetes.operator.metrics.MetricManager;
import io.fabric8.kubernetes.client.KubernetesClient;
+import io.fabric8.kubernetes.client.KubernetesClientException;
import io.fabric8.kubernetes.client.server.mock.EnableKubernetesMockClient;
import io.fabric8.kubernetes.client.server.mock.KubernetesMockServer;
import org.junit.jupiter.api.Test;
-import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
/** Test for {@link StatusRecorder}. */
@EnableKubernetesMockClient(crud = true)
@@ -47,17 +49,34 @@ public class StatusRecorderTest {
var lastRequest = mockServer.getLastRequest();
helper.patchAndCacheStatus(deployment, kubernetesClient);
- assertTrue(mockServer.getLastRequest() != lastRequest);
+ assertThat(lastRequest).isNotSameAs(mockServer.getLastRequest());
lastRequest = mockServer.getLastRequest();
deployment.getStatus().getReconciliationStatus().setState(ReconciliationState.ROLLING_BACK);
helper.patchAndCacheStatus(deployment, kubernetesClient);
// We intentionally compare references
- assertTrue(mockServer.getLastRequest() != lastRequest);
+ assertThat(lastRequest).isNotSameAs(mockServer.getLastRequest());
lastRequest = mockServer.getLastRequest();
// No update
helper.patchAndCacheStatus(deployment, kubernetesClient);
- assertTrue(mockServer.getLastRequest() == lastRequest);
+ assertThat(lastRequest).isSameAs(mockServer.getLastRequest());
+ }
+
+ @Test
+ public void testNullLatestResource() {
+ var statusRecorder =
+ new StatusRecorder<FlinkDeployment, FlinkDeploymentStatus>(
+ new MetricManager<>(), (e, s) -> {});
+
+ var resource = TestUtils.buildApplicationCluster();
+ var cause = new KubernetesClientException("dummy");
+ assertThatThrownBy(
+ () ->
+ statusRecorder.handleLockingError(
+ resource, null, kubernetesClient, 0,
cause))
+ .isInstanceOf(KubernetesClientException.class)
+ .hasMessage("Failed to retrieve latest resource")
+ .hasCause(cause);
}
}