Repository navigation
orchestrator-kubernetes: Fix replica metrics fetches failing on idle connections #39696
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,101 @@ | ||
| // Copyright Materialize, Inc. and contributors. All rights reserved. | ||
| // | ||
| // Use of this software is governed by the Business Source License | ||
| // included in the LICENSE file. | ||
| // | ||
| // As of the Change Date specified in that file, in accordance with | ||
| // the Business Source License, use of this software will be governed | ||
| // by the Apache License, Version 2.0. | ||
|
|
||
| //! Prometheus metrics for the Kubernetes orchestrator. | ||
|
|
||
| use std::error::Error; | ||
| use std::io; | ||
|
|
||
| use kube::error::Error as K8sError; | ||
| use mz_ore::metric; | ||
| use mz_ore::metrics::MetricsRegistry; | ||
| use mz_ore::metrics::raw::IntCounterVec; | ||
|
|
||
| /// Metrics for a [`KubernetesOrchestrator`](crate::KubernetesOrchestrator). | ||
| #[derive(Debug, Clone)] | ||
| pub(crate) struct OrchestratorMetrics { | ||
| process_metrics_fetch_failures: IntCounterVec, | ||
| } | ||
|
|
||
| /// The step of a per-process metrics fetch that failed. | ||
| #[derive(Debug, Clone, Copy)] | ||
| pub(crate) enum FetchStep { | ||
| /// Reading the pod's `PodMetrics` from the Kubernetes metrics API. | ||
| PodMetrics, | ||
| /// Reading usage metrics from the clusterd process. | ||
| ClusterdUsage, | ||
| } | ||
|
|
||
| impl FetchStep { | ||
| fn as_str(&self) -> &'static str { | ||
| match self { | ||
| FetchStep::PodMetrics => "pod_metrics", | ||
| FetchStep::ClusterdUsage => "clusterd_usage", | ||
| } | ||
| } | ||
| } | ||
|
|
||
| impl OrchestratorMetrics { | ||
| pub(crate) fn register_into(registry: &MetricsRegistry) -> Self { | ||
| Self { | ||
| process_metrics_fetch_failures: registry.register(metric!( | ||
| name: "mz_orchestrator_kubernetes_process_metrics_fetch_failures_total", | ||
| help: "The number of failed per-process metrics fetches, by fetch step and error kind.", | ||
| var_labels: ["step", "kind"], | ||
| )), | ||
| } | ||
| } | ||
|
|
||
| /// Records a failed fetch whose error has the given `kind`. | ||
| pub(crate) fn record_fetch_failure(&self, step: FetchStep, kind: &str) { | ||
| self.process_metrics_fetch_failures | ||
| .with_label_values(&[step.as_str(), kind]) | ||
| .inc(); | ||
| } | ||
|
|
||
| /// Records a failed fetch, classifying `error` by [`error_kind`]. | ||
| pub(crate) fn record_fetch_error(&self, step: FetchStep, error: &(dyn Error + 'static)) { | ||
| self.record_fetch_failure(step, error_kind(error)); | ||
| } | ||
| } | ||
|
|
||
| /// Classifies an error by walking its source chain. | ||
| /// | ||
| /// Returns `timeout` if any error in the chain is a timeout. Otherwise returns the kind of the | ||
| /// outermost error that is a Kubernetes client or HTTP client error: `api` for an error status | ||
| /// from the API server, `decode` for a response that failed to deserialize, `transport` for a | ||
| /// connection or protocol error. Returns `other` for anything else. | ||
| pub(crate) fn error_kind(error: &(dyn Error + 'static)) -> &'static str { | ||
| let mut kind = None; | ||
| let mut next = Some(error); | ||
| while let Some(err) = next { | ||
| next = err.source(); | ||
| if let Some(e) = err.downcast_ref::<io::Error>() { | ||
| if e.kind() == io::ErrorKind::TimedOut { | ||
| return "timeout"; | ||
| } | ||
| } else if let Some(e) = err.downcast_ref::<reqwest::Error>() { | ||
| if e.is_timeout() { | ||
| return "timeout"; | ||
| } | ||
| kind.get_or_insert_with(|| if e.is_decode() { "decode" } else { "transport" }); | ||
| } else if let Some(e) = err.downcast_ref::<K8sError>() { | ||
| kind.get_or_insert(match e { | ||
| K8sError::Api(_) => "api", | ||
| K8sError::SerdeError(_) => "decode", | ||
| K8sError::HyperError(_) | K8sError::Service(_) => "transport", | ||
| _ => "other", | ||
| }); | ||
| } | ||
| } | ||
| kind.unwrap_or("other") | ||
| } | ||
|
|
||
| #[cfg(test)] | ||
| mod tests; |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,52 @@ | ||
| // Copyright Materialize, Inc. and contributors. All rights reserved. | ||
| // | ||
| // Use of this software is governed by the Business Source License | ||
| // included in the LICENSE file. | ||
| // | ||
| // As of the Change Date specified in that file, in accordance with | ||
| // the Business Source License, use of this software will be governed | ||
| // by the Apache License, Version 2.0. | ||
|
|
||
| use std::io; | ||
|
|
||
| use anyhow::Context; | ||
| use kube::error::Error as K8sError; | ||
|
|
||
| use super::error_kind; | ||
|
|
||
| fn service_error(kind: io::ErrorKind) -> K8sError { | ||
| K8sError::Service(Box::new(io::Error::from(kind))) | ||
| } | ||
|
|
||
| #[mz_ore::test] | ||
| fn error_kind_finds_nested_timeout() { | ||
| assert_eq!( | ||
| error_kind(&service_error(io::ErrorKind::TimedOut)), | ||
| "timeout" | ||
| ); | ||
|
|
||
| let wrapped: anyhow::Result<()> = | ||
| Err(service_error(io::ErrorKind::TimedOut)).context("failed to get service"); | ||
| assert_eq!(error_kind(wrapped.unwrap_err().as_ref()), "timeout"); | ||
| } | ||
|
|
||
| #[mz_ore::test] | ||
| fn error_kind_classifies_kube_errors() { | ||
| assert_eq!( | ||
| error_kind(&service_error(io::ErrorKind::ConnectionReset)), | ||
| "transport" | ||
| ); | ||
|
|
||
| let serde_err = serde_json::from_str::<u64>("x").unwrap_err(); | ||
| assert_eq!(error_kind(&K8sError::SerdeError(serde_err)), "decode"); | ||
|
|
||
| let wrapped: anyhow::Result<()> = | ||
| Err(service_error(io::ErrorKind::ConnectionReset)).context("failed to get service"); | ||
| assert_eq!(error_kind(wrapped.unwrap_err().as_ref()), "transport"); | ||
| } | ||
|
|
||
| #[mz_ore::test] | ||
| fn error_kind_defaults_to_other() { | ||
| let err = anyhow::anyhow!("internal-http port missing in service spec"); | ||
| assert_eq!(error_kind(err.as_ref()), "other"); | ||
| } |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -38,7 +38,18 @@ pub async fn create_client(context: String) -> Result<(Client, String), anyhow:: | |
| }; | ||
|
|
||
| kubeconfig.connect_timeout = Some(Duration::from_secs(10)); | ||
| kubeconfig.read_timeout = Some(Duration::from_secs(60)); | ||
| // The read timeout bounds how long a hung Kubernetes call can block an orchestrator worker, | ||
| // which handles one command at a time. | ||
| // | ||
| // NOTE: The read timeout must exceed the connection pool's 90 s idle expiry. hyper-timeout | ||
| // starts the read timer while a pooled connection sits idle and does not reset it when a | ||
| // request is written. A request on a connection that idled for close to `read_timeout` | ||
| // therefore times out before the response arrives and fails with | ||
| // `client error (SendRequest)`. | ||
| // | ||
| // TODO(CPU-306): Return to 60 s once kube-client enables hyper-timeout's | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. kube-client won't enable it by default: the upstream change (kube-rs#2111) exposes it as an opt-in |
||
| // `reset_reader_on_write`. | ||
| kubeconfig.read_timeout = Some(Duration::from_secs(120)); | ||
| kubeconfig.write_timeout = Some(Duration::from_secs(60)); | ||
|
|
||
| let namespace = kubeconfig.default_namespace.clone(); | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Nit: exceeding 90 s isn't quite sufficient. Pool idle time counts against the next response for any read timeout, so at 120 s a request on a connection that idled ~89 s has ~31 s left. That's harmless for these millisecond-scale requests, but the actual constraint is "90 s plus the slowest response".