Skip to content

[SPARK-59928] Fix kubernetes.client.http.response.latency.nanos always recording near zero - #918

Open
ykisana wants to merge 1 commit into
apache:mainfrom
ykisana:SPARK-59928
Open

ykisana wants to merge 1 commit into
apache:mainfrom
ykisana:SPARK-59928

Conversation

@ykisana

@ykisana ykisana commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

What changes were proposed in this pull request?

Measure the response latency from a start time carried with the request, instead of taking two back-to-back System.nanoTime() calls in after(...). consumer(...) wraps the response body consumer in a TimedConsumer that carries the start time and the interceptor which created it (owner). after(...) finds it with AsyncBody.Consumer#unwrap(...), skipping any TimedConsumer of another KubernetesMetricsInterceptor on the same client, so that each interceptor uses its own start time.

The start time is restarted at each response. A request re-sent by fabric8 after another interceptor's afterFailure(...) returns true, e.g. TokenRefreshInterceptor on a 401, reuses the same consumer, so it is recorded twice, and the two values add up to the time from sending it until the re-sent response arrives, including the token refresh.

WebSocket upgrade responses (101) do not go through consumer(...), so after(...) finds no TimedConsumer and records no latency for them. They are still counted in kubernetes.client.http.response.

TimedConsumer#unwrap(...) delegates to the wrapped consumer, so fabric8's HttpLoggingInterceptor still works at the TRACE level.

Also:

  • Documents the latency in docs/configuration.md and adds an item to docs/migration_guide.md.
  • Corrects the Javadoc of after(...) (WebSocket upgrades and re-sent requests) and of afterConnectionFailure(...) (invoked only while the request has a retry left, even if it is not retried).
  • Moves the meterCount test helper to TestUtils, shared by KubernetesMetricsInterceptorTest and KubernetesClientFactoryTest.

Why are the changes needed?

Since SPARK-53647 / SPARK-53648 (0.5.0), kubernetes.client.http.response.latency.nanos has recorded nearly zero for every response, whatever the real latency, because both timestamps were taken in after(...).

Does this PR introduce any user-facing change?

Yes. kubernetes.client.http.response.latency.nanos changes in these ways; its name and type are unchanged:

  • It records the real response latency instead of nearly zero. It does not include the time to read the response body.
  • A request re-sent after a 401 is recorded twice, and the two values add up to the time from sending it until the re-sent response arrives, including the token refresh.
  • WebSocket upgrade responses (101), e.g. of a watch, are no longer recorded, so its count can be lower than that of kubernetes.client.http.response.

How was this patch tested?

Added tests to KubernetesMetricsInterceptorTest:

  • testResponseLatency: delays a mock server response by 50 ms and asserts that the recorded latency is at least 50 ms and below 30 s.
  • testResponseLatencyOfResentRequest: returns a 401 delayed by 200 ms and then a 200, which fabric8's TokenRefreshInterceptor re-sends. Asserts two latency samples, one of at least 200 ms and one below 200 ms, i.e. the re-sent response is not measured from the first attempt.
  • testResponseLatencyOfWebSocketUpgrade: starts an informer and asserts that the 101 upgrade response is counted in http.response but not in the latency histogram.
  • testResponseLatencyOfTwoInterceptors: registers two KubernetesMetricsInterceptors on one client, delays the response by 50 ms, and asserts that each records a latency of at least 50 ms, i.e. neither restarts the clock of the other.
  • testConsumerUnwrapsDelegate: asserts that TimedConsumer#unwrap delegates to the wrapped consumer.

Confirmed that testResponseLatencyOfResentRequest fails when the start time is not restarted at each response. Ran KubernetesMetricsInterceptorTest 10 times in a row without a failure.

Was this patch authored or co-authored using generative AI tooling?

Yes.

Generated-by: Claude Opus 5.5

@dongjoon-hyun

Copy link
Copy Markdown
Member

Thank you for making a PR, @ykisana .

@dongjoon-hyun dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you for making a PR, @ykisana. Moving the start time into consumer(...) fixes the near-zero values for regular HTTP requests, and testResponseLatency passes locally.

I left inline comments. Here is the summary.

  1. Build failure: The reordered import block violates the Spotless importOrder, so ./gradlew spotlessCheck (and ./gradlew build in CI) fails. Please run ./gradlew spotlessApply.
  2. Histogram population: fabric8 7.8.0 calls after(...) without consumer(...) for WebSocket upgrades (e.g. informer watches), and calls after(...) twice for a request re-sent after afterFailure(...) (e.g. 401 with TokenRefreshInterceptor). Those responses get no latency sample, so the count of http.response.latency.nanos no longer matches http.response.
  3. No-op cleanup: afterConnectionFailure(...) receives the original request, not the rebuilt one used as the key, so remove(request) never removes anything.
  4. Javadoc and tests: The new Javadoc implies that consumer(...) and after(...) are always paired, and the tests do not cover the paths above.
  5. Migration guide: Since this metric is enabled by default in 1.0.0, please add an item to docs/migration_guide.md.
  6. PR description: Please use the phrase Generated-by: Claude Opus 5.5 in the last section, as the PR template asks.
  7. Optional: fabric8 already passes the consumer returned by consumer(...) to after(...), so a delegating consumer which carries the start time could replace the global synchronizedMap(WeakHashMap).

@ykisana
ykisana force-pushed the SPARK-59928 branch 2 times, most recently from c567d5f to 6190979 Compare October 2, 2026 05:54
dongjoon-hyun added a commit that referenced this pull request Oct 2, 2026
…imers

### What changes were proposed in this pull request?

This PR fixes the Prometheus summaries of Dropwizard histograms and timers in `PrometheusPullModelHandler`.

- Track the sum of all recorded values in the new `SummingHistogram` and `SummingTimer`, which the operator now uses, and export `_sum` only for them.
- Export the median and the 99.9th percentile as the `0.5` and `0.999` quantiles of non-`nanos` histograms, instead of the mean and the 99th percentile.

### Why are the changes needed?

`_sum` was the mean of the reservoir multiplied by the count. Since the mean covers only roughly the last 5 minutes, `_sum` could decrease between scrapes, which breaks `rate()`. This becomes visible once [#918](#918) records real latencies.

### Does this PR introduce _any_ user-facing change?

Yes. The values of `_sum` and of the quantiles above change compared to [1.0.0](https://github.com/apache/spark-kubernetes-operator/releases/tag/1.0.0) (2026-07-26), while metric names and types are unchanged. The migration guide is updated.

### How was this patch tested?

Pass the CIs with the newly added test cases.

### Was this patch authored or co-authored using generative AI tooling?

Generated-by: Claude Opus 5.5

Closes #922 from dongjoon-hyun/SPARK-59935.

Authored-by: Dongjoon Hyun <dongjoon@apache.org>
Signed-off-by: Dongjoon Hyun <dongjoon@apache.org>

@dongjoon-hyun dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you for updating the PR, @ykisana. The delegating TimedConsumer fixes the near-zero values of regular requests, and the linters, javadoc and the tests pass locally.

I left inline comments. Here is the summary.

  1. Re-sent request: fabric8 re-sends a request after afterFailure(...) with the same consumer. So the latency of the re-sent response is measured from the start of the first attempt, and the first exchange is counted twice, e.g. 401 -> token refresh -> 200. A 5xx or 429 retry gets a fresh consumer instead.
  2. 401 latency: after(...) of a failed response runs only after the synchronous afterFailure(...) work. So a 401 sample can include the token refresh, e.g. a kubeconfig exec credential plugin.
  3. Rebase: This PR conflicts with main in docs/migration_guide.md because of SPARK-59935. SPARK-59935 is also needed for a correct Prometheus _sum of this histogram.
  4. Flaky test: testResponseLatencyOfWebSocketUpgrade awaits only http.response, so it can fail with an NPE at http.response.101.
  5. Javadoc and docs: The re-send sentence of after(...) does not hold for a WebSocket upgrade. after(...) gets the outermost consumer of the chain, not necessarily the one returned by consumer(...). The latency ends when the response headers arrive. 1.0 recorded less than a microsecond.
  6. Test coverage: No test covers the delegation in TimedConsumer#unwrap, which HttpLoggingInterceptor needs at the TRACE level.
  7. PR description: Please add Generated-by: Claude Opus 5.5 as the PR template asks. Please also update How was this patch tested?, which still describes a 500 ms delay while the test uses 50 ms and does not mention the two other tests, and Does this PR introduce any user-facing change?, which does not mention that the count no longer includes WebSocket upgrades.
  8. Nits: an unneeded await(), the nullable Long parameter, the anonymous 403 re-send interceptor, and the metric lookups in the tests.

Comment thread docs/migration_guide.md Outdated
Comment thread docs/migration_guide.md Outdated
@ykisana
ykisana force-pushed the SPARK-59928 branch 3 times, most recently from c377607 to 0049941 Compare October 3, 2026 05:52
@ykisana

ykisana commented Oct 3, 2026

Copy link
Copy Markdown
Contributor Author

Thank you for updating the PR, @ykisana. The delegating TimedConsumer fixes the near-zero values of regular requests, and the linters, javadoc and the tests pass locally.

I left inline comments. Here is the summary.

1. **Re-sent request**: fabric8 re-sends a request after `afterFailure(...)` with the same consumer. So the latency of the re-sent response is measured from the start of the first attempt, and the first exchange is counted twice, e.g. `401` -> token refresh -> `200`. A `5xx` or `429` retry gets a fresh consumer instead.

2. **`401` latency**: `after(...)` of a failed response runs only after the synchronous `afterFailure(...)` work. So a `401` sample can include the token refresh, e.g. a kubeconfig exec credential plugin.

3. **Rebase**: This PR conflicts with `main` in `docs/migration_guide.md` because of [SPARK-59935](https://issues.apache.org/jira/browse/SPARK-59935). [SPARK-59935](https://issues.apache.org/jira/browse/SPARK-59935) is also needed for a correct Prometheus `_sum` of this histogram.

4. **Flaky test**: `testResponseLatencyOfWebSocketUpgrade` awaits only `http.response`, so it can fail with an NPE at `http.response.101`.

5. **Javadoc and docs**: The re-send sentence of `after(...)` does not hold for a WebSocket upgrade. `after(...)` gets the outermost consumer of the chain, not necessarily the one returned by `consumer(...)`. The latency ends when the response headers arrive. 1.0 recorded less than a microsecond.

6. **Test coverage**: No test covers the delegation in `TimedConsumer#unwrap`, which `HttpLoggingInterceptor` needs at the `TRACE` level.

7. **PR description**: Please add `Generated-by: Claude Opus 5.5` as the PR template asks. Please also update `How was this patch tested?`, which still describes a 500 ms delay while the test uses 50 ms and does not mention the two other tests, and `Does this PR introduce any user-facing change?`, which does not mention that the count no longer includes WebSocket upgrades.

8. **Nits**: an unneeded `await()`, the nullable `Long` parameter, the anonymous `403` re-send interceptor, and the metric lookups in the tests.

Addressed these, thanks for looking into it so much.
Opus was used to address all these, updated the PR desc.

@dongjoon-hyun dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you for addressing all the comments, @ykisana. Restarting the clock at each response fixes the re-sent request. The linters, javadoc and the tests pass locally. I also confirmed that testResponseLatencyOfResentRequest fails when getAndSet(now) is replaced with get(), and that testConsumerUnwrapsDelegate fails without the delegation in TimedConsumer#unwrap.

I left a few more inline comments. Here is the summary.

  1. Asynchronous token refresh: With an OIDC auth provider, fabric8 refreshes the token asynchronously, so the refresh time goes into the latency of the re-sent request, not of the 401. The new sentence in docs/configuration.md holds only for a synchronous refresh.
  2. PR description:
    • Does this PR introduce _any_ user-facing change? still does not mention two changes: WebSocket upgrades are no longer recorded, so the count can be lower than that of kubernetes.client.http.response, and the latency ends when the response headers arrive.
    • Why are the changes needed? still says only a few microseconds per response, while the migration guide says nearly zero.
    • How was this patch tested? lists testTimedConsumerUnwrapsDelegate, but the test is testConsumerUnwrapsDelegate.
    • What changes were proposed in this pull request? does not mention that the start time is restarted at each response. Also, from when the request is sent does not hold for a re-sent request.
  3. Two metrics interceptors (corner case): getAndSet(...) restarts a clock which every KubernetesMetricsInterceptor on a client shares. So a second instance, e.g. of a subclass, records nearly zero.
  4. Nits: The 401 delay in the test relies on null being serialized into a non-empty body, and the histogram lookups in the tests could be simpler.

Comment thread docs/configuration.md Outdated
`PodsReady` condition, is not counted there when it gets no response.

`kubernetes.client.http.response.latency.nanos` ends when the response headers arrive, so it does
not include reading the body. A request re-sent after a `401` is measured from the `401`, whose

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This holds when the token is refreshed synchronously, e.g. with a service account token file, an exec credential plugin or an OAuthTokenProvider. But with an OIDC auth provider, TokenRefreshInterceptor#afterFailure returns the pending future of OpenIDConnectionUtils#resolveOIDCTokenFromAuthConfig, which calls the IdP with sendAsync. Then fabric8 runs after(...) of the 401 before the refresh, and re-sends the request only when the refresh completes. So the refresh time goes into the latency of the re-sent request instead. In my test with a 300 ms refresh, a 401 returned after 100 ms was recorded as 105 ms, and the re-sent 200 returned after 50 ms as ~360 ms. With a synchronous refresh of the same length, they were 410 ms and 54 ms.

The sum of the two values is the same either way. How about describing that instead? e.g.

A request re-sent after a 401 is recorded twice, and the two values add up to the time from sending it until the re-sent response arrives, including the token refresh.

Could you also align the Javadoc of after(...) (L132) and TimedConsumer (L262-263) with it? The same goes for from sending each HTTP request in the migration guide.

.get()
.delay(200)
.withPath(CONFIG_MAP_PATH)
.andReturn(HTTP_UNAUTHORIZED, null)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: The .delay(200) works here only because null is serialized into the 4-byte body null. fabric8 mockwebserver 7.8.0 delays only a non-empty body (HttpServerRequestHandler). With an empty body, the delay would be silently dropped, and the assertion at L193 would fail with only expected: <true> but was: <false>. How about returning an explicit body, e.g. a Status with the code 401? Could you also add messages with the snapshot values to the timing assertions?

client.resource(configMap).get();

Map<String, Metric> map = metricsInterceptor.metricRegistry().getMetrics();
Histogram latency = (Histogram) map.get("http.response.latency.nanos");

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: Meters now go through meterCount(...), but the three tests still keep a Map<String, Metric> only for the histogram cast, and the WebSocket test mixes both styles in one untilAsserted(...). How about metricsInterceptor.metricRegistry().getHistograms().get("http.response.latency.nanos"), or a helper next to meterCount(...)? FYI, KubernetesClientFactoryTest#meterCount has the same name and signature but creates a missing meter, so a misspelled name returns 0 there and fails here.

@ykisana
ykisana force-pushed the SPARK-59928 branch 2 times, most recently from 0f38e3d to dcd39df Compare October 5, 2026 23:18
@ykisana

ykisana commented Oct 5, 2026

Copy link
Copy Markdown
Contributor Author

@dongjoon-hyun address the comments

Please let me know if any more issues

@ykisana

ykisana commented Oct 7, 2026

Copy link
Copy Markdown
Contributor Author

Updated to fix merge conflict

@dongjoon-hyun

Copy link
Copy Markdown
Member

Thank you for updating. I'm re-reviewing now, @ykisana .

@ykisana

ykisana commented Oct 7, 2026

Copy link
Copy Markdown
Contributor Author

Removed JIRA reference here too as per: #944

@dongjoon-hyun dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you for updating the PR, @ykisana. The owner-based lookup fixes the two-interceptor case. The linters, javadoc and the tests pass locally, and the four new latency tests passed 30 repeated runs each. I also confirmed with fabric8 7.8.0 and Vert.x 4.5.33 that the two values of a re-sent request add up to the total time, and that HttpLoggingInterceptor still finds its consumer at the TRACE level.

I left inline comments. Here is the summary.

  1. 401 latency in the docs: With a synchronous token refresh, the 401 value ends after the refresh, not when its headers arrive. Only the sum holds.
  2. Rejected WebSocket upgrades: With the default interceptors, they are always counted five times, also in the response code meters.
  3. Javadoc:
    • afterConnectionFailure(...): fabric8 invokes it for the last attempt of a request timeout while retries are left.
    • after(...): {@link #afterFailure(...)} points to this class's method, which never re-sends.
  4. PR description: testResponseLatencyOfTwoInterceptors, the owner lookup and the docs changes are missing, and the per-value re-send sentence holds only for an asynchronous token refresh.
  5. Nits: pollDelay(Duration.ZERO) in the WebSocket test, the meterCount helper, and an optional upstream fabric8 issue for the re-send path.
  6. Pre-existing issues (FYI, not from this PR): the metric names in the client metrics table, an IndexOutOfBoundsException for a status code outside 100-599, and the .nanos name with values in seconds when the Prometheus name sanitizing is disabled.

Comment thread docs/configuration.md Outdated
Comment thread docs/configuration.md Outdated
Comment thread docs/configuration.md
Comment thread docs/migration_guide.md
@ykisana

ykisana commented Oct 7, 2026 •

Copy link
Copy Markdown
Contributor Author

@dongjoon-hyun

Addressed in this PR

  • docs/configuration.md and docs/migration_guide.md: dropped "ends when the response headers arrive" / "until their headers arrive" and now say it does not include the time to read the response body. The 401 sentence only describes the sum.
  • docs/configuration.md: a rejected WebSocket upgrade is now described as counted several times in kubernetes.client.http.response and the response code meters.
  • after(...) Javadoc: the re-send now refers to another interceptor's afterFailure(...) returning true, e.g. TokenRefreshInterceptor on a 401, instead of linking to this class's afterFailure(...).
  • afterConnectionFailure(...) Javadoc: reused the docs wording, "only while the request has a retry left, even if it does not retry it, e.g. on a request timeout."
  • testResponseLatencyOfWebSocketUpgrade: added pollDelay(Duration.ZERO).
  • Moved meterCount to TestUtils, shared by KubernetesMetricsInterceptorTest and KubernetesClientFactoryTest.
  • PR description: added testResponseLatencyOfTwoInterceptors, the owner lookup, and the docs, Javadoc and helper changes. Also described the 401 re-send as the sum of the two values.

To be addressed separately (new issues/PRs)

  • Client metrics table in docs/configuration.md: fix the names of kubernetes.client.failed and kubernetes.client.1xx-5xx, and the description of failed, which is marked for every non-2xx response that reaches the interceptor.
  • IndexOutOfBoundsException for a status code outside 100-599 in the response code group meters: add a bounds check.
  • With spark.kubernetes.operator.metrics.sanitizePrometheusMetricsNameEnabled=false, the .nanos name is kept while the values are in seconds.
  • File an upstream fabric8 issue for the re-send path (after(...) called again without before(...)/consumer(...)), and the WebSocket upgrade paths (no after(...) for a re-sent upgrade, 1+N calls for a rejected one).

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants