Skip to content

Commit dcd39df

Browse files
committed
[SPARK-59928] Fix kubernetes.client.http.response.latency.nanos always recording near zero
1 parent 833f05f commit dcd39df

4 files changed

Lines changed: 218 additions & 7 deletions

File tree

‎docs/configuration.md‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -263,6 +263,12 @@ has a retry left for it, even if it does not retry it, e.g. on a request timeout
263263
a client which never retries, i.e. a Kubernetes event write or the recording of the Kueue
264264
`PodsReady` condition, is not counted there when it gets no response.
265265

266+
`kubernetes.client.http.response.latency.nanos` ends when the response headers arrive, so it does
267+
not include reading the body. A request re-sent after a `401` is recorded twice, and the two
268+
values add up to the time from sending it until the re-sent response arrives, including the token
269+
refresh. WebSocket upgrades, e.g. of a watch, are not recorded, and a rejected one can be counted
270+
more than once in `kubernetes.client.http.response`.
271+
266272
### Latency for State Transition
267273

268274
Spark Operator also measures the latency between each state transition for apps, in the format of

‎docs/migration_guide.md‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -116,6 +116,12 @@ affect building the project, CI, tests or examples.
116116
`NullPointerException` for most properties. The operator logs each invalid value once as a
117117
warning, instead of as an error with a stack trace whenever it reads the value. Use `true` or
118118
`false`, which behave the same in both versions ([SPARK-XXXXX](https://issues.apache.org/jira/browse/SPARK-XXXXX)).
119+
- Since 1.1.0, `kubernetes.client.http.response.latency.nanos` records the latency of HTTP
120+
responses until their headers arrive. 1.0 recorded nearly zero regardless of the latency. A
121+
request re-sent after a `401` is recorded twice, and the two values add up to the time from
122+
sending it until the re-sent response arrives, including the token refresh. It no longer records
123+
a WebSocket upgrade response, e.g. of a watch, so its count can be
124+
lower than that of `kubernetes.client.http.response` ([SPARK-59928](https://issues.apache.org/jira/browse/SPARK-59928)).
119125

120126
### SparkApplication
121127

‎spark-operator/src/main/java/org/apache/spark/k8s/operator/metrics/source/KubernetesMetricsInterceptor.java‎

Lines changed: 62 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
import java.util.Optional;
3030
import java.util.concurrent.CompletableFuture;
3131
import java.util.concurrent.ConcurrentHashMap;
32+
import java.util.concurrent.atomic.AtomicLong;
3233

3334
import com.codahale.metrics.Histogram;
3435
import com.codahale.metrics.Meter;
@@ -90,6 +91,25 @@ public KubernetesMetricsInterceptor() {
9091
}
9192
}
9293

94+
/**
95+
* Called before a request is sent, to allow wrapping the consumer of the response body.
96+
*
97+
* <p>Wraps the consumer to carry the request start time, which {@link #after(HttpRequest,
98+
* HttpResponse, AsyncBody.Consumer)} finds with {@link AsyncBody.Consumer#unwrap(Class)}. fabric8
99+
* passes to it the outermost consumer after all interceptors have wrapped it, so an interceptor
100+
* which wraps the consumer after this one must delegate {@code unwrap(...)}, or no latency is
101+
* recorded.
102+
*
103+
* @param consumer the consumer of the response body
104+
* @param request the request about to be sent
105+
* @return the consumer carrying the request start time
106+
*/
107+
@Override
108+
public AsyncBody.Consumer<List<ByteBuffer>> consumer(
109+
AsyncBody.Consumer<List<ByteBuffer>> consumer, HttpRequest request) {
110+
return new TimedConsumer(consumer, this, new AtomicLong(System.nanoTime()));
111+
}
112+
93113
/**
94114
* Called before a request to allow for the manipulation of the request
95115
*
@@ -102,21 +122,35 @@ public void before(BasicBuilder builder, HttpRequest request, RequestTags tags)
102122
}
103123

104124
/**
105-
* Called after a non-WebSocket HTTP response is received. The body might or might not be already
106-
* consumed.
125+
* Called after an HTTP response is received. The body might or might not be already consumed.
107126
*
108127
* <p>Should be used to analyze response codes and headers, original response shouldn't be
109128
* altered.
110129
*
130+
* <p>fabric8 also calls this for a WebSocket upgrade response with a {@code null} consumer, and
131+
* again for an HTTP request re-sent after {@link #afterFailure(BasicBuilder, HttpResponse,
132+
* RequestTags)}. The two values of a request re-sent this way add up to the time from sending it
133+
* until the re-sent response arrives.
134+
*
111135
* @param request the original request sent to the server.
112136
* @param response the response received from the server.
137+
* @param consumer the consumer of the response body, or {@code null} for a WebSocket upgrade.
113138
*/
114139
@Override
140+
@SuppressWarnings("PMD.CompareObjectsWithEquals")
115141
public void after(
116142
HttpRequest request,
117143
HttpResponse<?> response,
118144
AsyncBody.Consumer<List<ByteBuffer>> consumer) {
119-
updateResponseMetrics(response, System.nanoTime());
145+
TimedConsumer timedConsumer = consumer == null ? null : consumer.unwrap(TimedConsumer.class);
146+
while (timedConsumer != null && timedConsumer.owner() != this) {
147+
timedConsumer = timedConsumer.delegate().unwrap(TimedConsumer.class);
148+
}
149+
if (timedConsumer != null) {
150+
long now = System.nanoTime();
151+
responseLatency.update(now - timedConsumer.startTimeNanos().getAndSet(now));
152+
}
153+
updateResponseMetrics(response);
120154
}
121155

122156
/**
@@ -139,7 +173,8 @@ public CompletableFuture<Boolean> afterFailure(
139173
/**
140174
* Called after a connection attempt fails.
141175
*
142-
* <p>This method will be invoked on each failed connection attempt.
176+
* <p>fabric8 invokes this method only for a failed connection attempt which has a retry left, not
177+
* for the last attempt.
143178
*
144179
* @param request the HTTP request.
145180
* @param failure the Java exception that caused the failure.
@@ -183,11 +218,9 @@ private void updateRequestMetrics(HttpRequest request) {
183218
});
184219
}
185220

186-
private void updateResponseMetrics(HttpResponse response, long startTimeNanos) {
221+
private void updateResponseMetrics(HttpResponse response) {
187222
Objects.requireNonNull(response);
188-
final long latency = System.nanoTime() - startTimeNanos;
189223
responseRateMeter.mark();
190-
responseLatency.update(latency);
191224
getMeterByResponseCode(response.code()).mark();
192225
if (KUBERNETES_CLIENT_METRICS_GROUP_BY_RESPONSE_CODE_GROUP_ENABLED.getValue()) {
193226
responseCodeGroupMeters.get(response.code() / 100 - 1).mark();
@@ -229,4 +262,26 @@ public Optional<Pair<String, String>> parseNamespaceScopedResource(String path)
229262
return Optional.empty();
230263
}
231264
}
265+
266+
/**
267+
* Consumer of a response body which carries the start time of its request, restarted at each
268+
* response so that the values of a re-sent request add up to its latency. It also carries the
269+
* interceptor which created it, so that each interceptor on a client uses its own start time.
270+
*/
271+
private record TimedConsumer(
272+
AsyncBody.Consumer<List<ByteBuffer>> delegate,
273+
KubernetesMetricsInterceptor owner,
274+
AtomicLong startTimeNanos)
275+
implements AsyncBody.Consumer<List<ByteBuffer>> {
276+
@Override
277+
public void consume(List<ByteBuffer> value, AsyncBody asyncBody) throws Exception {
278+
delegate.consume(value, asyncBody);
279+
}
280+
281+
@Override
282+
public <U> U unwrap(Class<U> target) {
283+
U self = AsyncBody.Consumer.super.unwrap(target);
284+
return self != null ? self : delegate.unwrap(target);
285+
}
286+
}
232287
}

‎spark-operator/src/test/java/org/apache/spark/k8s/operator/metrics/source/KubernetesMetricsInterceptorTest.java‎

Lines changed: 144 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,19 +19,29 @@
1919

2020
package org.apache.spark.k8s.operator.metrics.source;
2121

22+
import static java.net.HttpURLConnection.HTTP_UNAUTHORIZED;
23+
import static org.awaitility.Awaitility.await;
2224
import static org.junit.jupiter.api.Assertions.assertThrows;
2325

26+
import java.nio.ByteBuffer;
2427
import java.util.Arrays;
2528
import java.util.HashMap;
2629
import java.util.List;
2730
import java.util.Map;
2831

32+
import com.codahale.metrics.Histogram;
2933
import com.codahale.metrics.Meter;
3034
import com.codahale.metrics.Metric;
35+
import com.codahale.metrics.Snapshot;
3136
import io.fabric8.kubernetes.api.model.ConfigMap;
3237
import io.fabric8.kubernetes.api.model.ObjectMeta;
38+
import io.fabric8.kubernetes.api.model.StatusBuilder;
39+
import io.fabric8.kubernetes.client.Config;
40+
import io.fabric8.kubernetes.client.ConfigBuilder;
3341
import io.fabric8.kubernetes.client.KubernetesClient;
42+
import io.fabric8.kubernetes.client.http.AsyncBody;
3443
import io.fabric8.kubernetes.client.http.Interceptor;
44+
import io.fabric8.kubernetes.client.informers.SharedIndexInformer;
3545
import io.fabric8.kubernetes.client.server.mock.EnableKubernetesMockClient;
3646
import io.fabric8.kubernetes.client.server.mock.KubernetesMockServer;
3747
import org.junit.jupiter.api.AfterEach;
@@ -50,6 +60,9 @@
5060
@TestMethodOrder(MethodOrderer.OrderAnnotation.class)
5161
class KubernetesMetricsInterceptorTest {
5262

63+
private static final String CONFIG_MAP_PATH =
64+
"/api/v1/namespaces/spark-system/configmaps/spark-job-operator-configuration";
65+
5366
private KubernetesMockServer mockServer;
5467
private KubernetesClient kubernetesClient;
5568

@@ -125,6 +138,137 @@ void testWhenKubernetesServerNotWorking() {
125138
}
126139
}
127140

141+
@Test
142+
@Order(3)
143+
void testResponseLatency() {
144+
KubernetesMetricsInterceptor metricsInterceptor = new KubernetesMetricsInterceptor();
145+
List<Interceptor> interceptors = List.of(metricsInterceptor);
146+
try (KubernetesClient client =
147+
KubernetesClientFactory.buildKubernetesClient(
148+
interceptors, kubernetesClient.getConfiguration())) {
149+
ConfigMap configMap = createConfigMap();
150+
mockServer
151+
.expect()
152+
.get()
153+
.delay(50)
154+
.withPath(CONFIG_MAP_PATH)
155+
.andReturn(200, configMap)
156+
.once();
157+
client.resource(configMap).get();
158+
159+
Histogram latency = latency(metricsInterceptor);
160+
Assertions.assertEquals(1, latency.getCount());
161+
Snapshot snapshot = latency.getSnapshot();
162+
Assertions.assertTrue(snapshot.getMax() >= 50_000_000L, () -> "max: " + snapshot.getMax());
163+
Assertions.assertTrue(
164+
snapshot.getMax() < 30_000_000_000L, () -> "max: " + snapshot.getMax());
165+
}
166+
}
167+
168+
@Test
169+
@Order(4)
170+
void testResponseLatencyOfResentRequest() {
171+
KubernetesMetricsInterceptor metricsInterceptor = new KubernetesMetricsInterceptor();
172+
List<Interceptor> interceptors = List.of(metricsInterceptor);
173+
Config config =
174+
new ConfigBuilder(kubernetesClient.getConfiguration())
175+
.withOauthTokenProvider(() -> "token")
176+
.build();
177+
try (KubernetesClient client =
178+
KubernetesClientFactory.buildKubernetesClient(interceptors, config)) {
179+
ConfigMap configMap = createConfigMap();
180+
mockServer
181+
.expect()
182+
.get()
183+
.delay(200)
184+
.withPath(CONFIG_MAP_PATH)
185+
.andReturn(HTTP_UNAUTHORIZED, new StatusBuilder().withCode(HTTP_UNAUTHORIZED).build())
186+
.once();
187+
mockServer.expect().get().withPath(CONFIG_MAP_PATH).andReturn(200, configMap).once();
188+
client.resource(configMap).get();
189+
190+
Assertions.assertEquals(2, meterCount(metricsInterceptor, "http.response"));
191+
Assertions.assertEquals(1, meterCount(metricsInterceptor, "http.response.401"));
192+
Assertions.assertEquals(0, meterCount(metricsInterceptor, "failed"));
193+
Histogram latency = latency(metricsInterceptor);
194+
Assertions.assertEquals(2, latency.getCount());
195+
Snapshot snapshot = latency.getSnapshot();
196+
Assertions.assertTrue(snapshot.getMax() >= 200_000_000L, () -> "max: " + snapshot.getMax());
197+
Assertions.assertTrue(snapshot.getMin() < 200_000_000L, () -> "min: " + snapshot.getMin());
198+
}
199+
}
200+
201+
@Test
202+
@Order(5)
203+
void testResponseLatencyOfWebSocketUpgrade() {
204+
KubernetesMetricsInterceptor metricsInterceptor = new KubernetesMetricsInterceptor();
205+
List<Interceptor> interceptors = List.of(metricsInterceptor);
206+
try (KubernetesClient client =
207+
KubernetesClientFactory.buildKubernetesClient(
208+
interceptors, kubernetesClient.getConfiguration());
209+
SharedIndexInformer<ConfigMap> informer =
210+
client.configMaps().inNamespace("spark-system").inform()) {
211+
await()
212+
.untilAsserted(
213+
() -> {
214+
Assertions.assertEquals(2, meterCount(metricsInterceptor, "http.response"));
215+
Assertions.assertEquals(1, meterCount(metricsInterceptor, "http.response.101"));
216+
Assertions.assertEquals(1, latency(metricsInterceptor).getCount());
217+
});
218+
Assertions.assertTrue(informer.isWatching());
219+
}
220+
}
221+
222+
@Test
223+
@Order(6)
224+
void testResponseLatencyOfTwoInterceptors() {
225+
KubernetesMetricsInterceptor first = new KubernetesMetricsInterceptor();
226+
KubernetesMetricsInterceptor second = new KubernetesMetricsInterceptor() {};
227+
List<Interceptor> interceptors = List.of(first, second);
228+
try (KubernetesClient client =
229+
KubernetesClientFactory.buildKubernetesClient(
230+
interceptors, kubernetesClient.getConfiguration())) {
231+
ConfigMap configMap = createConfigMap();
232+
mockServer
233+
.expect()
234+
.get()
235+
.delay(50)
236+
.withPath(CONFIG_MAP_PATH)
237+
.andReturn(200, configMap)
238+
.once();
239+
client.resource(configMap).get();
240+
241+
for (KubernetesMetricsInterceptor interceptor : List.of(first, second)) {
242+
Histogram latency = latency(interceptor);
243+
Assertions.assertEquals(1, latency.getCount());
244+
Snapshot snapshot = latency.getSnapshot();
245+
Assertions.assertTrue(snapshot.getMin() >= 50_000_000L, () -> "min: " + snapshot.getMin());
246+
}
247+
}
248+
}
249+
250+
@Test
251+
@Order(7)
252+
void testConsumerUnwrapsDelegate() {
253+
AsyncBody.Consumer<List<ByteBuffer>> delegate = (value, asyncBody) -> {};
254+
AsyncBody.Consumer<List<ByteBuffer>> consumer =
255+
new KubernetesMetricsInterceptor().consumer(delegate, null);
256+
Assertions.assertSame(delegate, consumer.unwrap(delegate.getClass()));
257+
}
258+
259+
private static long meterCount(KubernetesMetricsInterceptor interceptor, String name) {
260+
Meter meter = interceptor.metricRegistry().getMeters().get(name);
261+
Assertions.assertNotNull(meter, name);
262+
return meter.getCount();
263+
}
264+
265+
private static Histogram latency(KubernetesMetricsInterceptor interceptor) {
266+
Histogram histogram =
267+
interceptor.metricRegistry().getHistograms().get("http.response.latency.nanos");
268+
Assertions.assertNotNull(histogram, "http.response.latency.nanos");
269+
return histogram;
270+
}
271+
128272
private static SparkApplication createSparkApplication() {
129273
ObjectMeta meta = new ObjectMeta();
130274
meta.setName("sample-spark-application");

0 commit comments

Comments
 (0)