Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions doc/user/data/metrics.yml
Original file line number Diff line number Diff line change
Expand Up @@ -1687,6 +1687,13 @@ metrics:
- case
source: src/compute-types/src/plan.rs
visibility: internal
- name: mz_orchestrator_kubernetes_process_metrics_fetch_failures_total
help: The number of failed per-process metrics fetches, by fetch step and error kind.
labels:
- kind
- step
source: src/orchestrator-kubernetes/src/metrics.rs
visibility: internal
- name: mz_otel_on_close
help: count of on_close events sent to otel
source: src/ore/src/tracing.rs
Expand Down
73 changes: 38 additions & 35 deletions src/environmentd/src/environmentd/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -824,41 +824,44 @@ fn run(mut args: Args) -> Result<(), anyhow::Error> {

let orchestrator = Arc::new(
runtime
.block_on(KubernetesOrchestrator::new(KubernetesOrchestratorConfig {
context: args.orchestrator_kubernetes_context.clone(),
scheduler_name: args.orchestrator_kubernetes_scheduler_name,
priority_class_name: args.orchestrator_kubernetes_priority_class_name,
service_annotations: args
.orchestrator_kubernetes_service_annotation
.into_iter()
.map(|l| (l.key, l.value))
.collect(),
service_labels: args
.orchestrator_kubernetes_service_label
.into_iter()
.map(|l| (l.key, l.value))
.collect(),
service_node_selector: args
.orchestrator_kubernetes_service_node_selector
.into_iter()
.map(|l| (l.key, l.value))
.collect(),
service_affinity: args.orchestrator_kubernetes_service_affinity,
service_tolerations: args.orchestrator_kubernetes_service_tolerations,
service_account: args.orchestrator_kubernetes_service_account,
image_pull_policy: args.orchestrator_kubernetes_image_pull_policy,
aws_external_id_prefix: args.aws_external_id_prefix.clone(),
coverage: args.orchestrator_kubernetes_coverage,
ephemeral_volume_storage_class: args
.orchestrator_kubernetes_ephemeral_volume_class
.clone(),
service_fs_group: args.orchestrator_kubernetes_service_fs_group.clone(),
name_prefix: args.orchestrator_kubernetes_name_prefix.clone(),
collect_pod_metrics: !args
.orchestrator_kubernetes_disable_pod_metrics_collection,
enable_prometheus_scrape_annotations: args
.orchestrator_kubernetes_enable_prometheus_scrape_annotations,
}))
.block_on(KubernetesOrchestrator::new(
KubernetesOrchestratorConfig {
context: args.orchestrator_kubernetes_context.clone(),
scheduler_name: args.orchestrator_kubernetes_scheduler_name,
priority_class_name: args.orchestrator_kubernetes_priority_class_name,
service_annotations: args
.orchestrator_kubernetes_service_annotation
.into_iter()
.map(|l| (l.key, l.value))
.collect(),
service_labels: args
.orchestrator_kubernetes_service_label
.into_iter()
.map(|l| (l.key, l.value))
.collect(),
service_node_selector: args
.orchestrator_kubernetes_service_node_selector
.into_iter()
.map(|l| (l.key, l.value))
.collect(),
service_affinity: args.orchestrator_kubernetes_service_affinity,
service_tolerations: args.orchestrator_kubernetes_service_tolerations,
service_account: args.orchestrator_kubernetes_service_account,
image_pull_policy: args.orchestrator_kubernetes_image_pull_policy,
aws_external_id_prefix: args.aws_external_id_prefix.clone(),
coverage: args.orchestrator_kubernetes_coverage,
ephemeral_volume_storage_class: args
.orchestrator_kubernetes_ephemeral_volume_class
.clone(),
service_fs_group: args.orchestrator_kubernetes_service_fs_group.clone(),
name_prefix: args.orchestrator_kubernetes_name_prefix.clone(),
collect_pod_metrics: !args
.orchestrator_kubernetes_disable_pod_metrics_collection,
enable_prometheus_scrape_annotations: args
.orchestrator_kubernetes_enable_prometheus_scrape_annotations,
},
&metrics_registry,
))
.context("creating kubernetes orchestrator")?,
);
let secrets_controller: Arc<dyn SecretsController> = match args.secrets_controller {
Expand Down
3 changes: 3 additions & 0 deletions src/orchestrator-kubernetes/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -32,5 +32,8 @@ sha2.workspace = true
tokio.workspace = true
tracing.workspace = true

[dev-dependencies]
mz-ore = { path = "../ore", default-features = false, features = ["test"] }

[features]
default = []
47 changes: 35 additions & 12 deletions src/orchestrator-kubernetes/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -52,14 +52,19 @@ use mz_orchestrator::{
ServiceProcessMetrics, ServiceStatus, recommended_k8s_labels, scheduling_config::*,
};
use mz_ore::cast::CastInto;
use mz_ore::error::ErrorExt;
use mz_ore::metrics::MetricsRegistry;
use mz_ore::retry::Retry;
use mz_ore::task::AbortOnDropHandle;
use serde::Deserialize;
use sha2::{Digest, Sha256};
use tokio::sync::{mpsc, oneshot};
use tracing::{error, info, warn};

use crate::metrics::{FetchStep, OrchestratorMetrics};

pub mod cloud_resource_controller;
mod metrics;
pub mod secrets;
pub mod util;

Expand Down Expand Up @@ -164,6 +169,7 @@ pub struct KubernetesOrchestrator {
vpc_endpoint_api: Api<VpcEndpoint>,
namespaces: Mutex<BTreeMap<String, Arc<dyn NamespacedOrchestrator>>>,
resource_reader: Arc<KubernetesResourceReader>,
metrics: OrchestratorMetrics,
}

impl fmt::Debug for KubernetesOrchestrator {
Expand All @@ -176,6 +182,7 @@ impl KubernetesOrchestrator {
/// Creates a new Kubernetes orchestrator from the provided configuration.
pub async fn new(
config: KubernetesOrchestratorConfig,
metrics_registry: &MetricsRegistry,
) -> Result<KubernetesOrchestrator, anyhow::Error> {
let (client, kubernetes_namespace) = util::create_client(config.context.clone()).await?;
let resource_reader =
Expand All @@ -188,6 +195,7 @@ impl KubernetesOrchestrator {
vpc_endpoint_api: Api::default_namespaced(client),
namespaces: Mutex::new(BTreeMap::new()),
resource_reader,
metrics: OrchestratorMetrics::register_into(metrics_registry),
})
}
}
Expand All @@ -206,6 +214,7 @@ impl Orchestrator for KubernetesOrchestrator {
command_rx,
name_prefix: self.config.name_prefix.clone().unwrap_or_default(),
collect_pod_metrics: self.config.collect_pod_metrics,
metrics: self.metrics.clone(),
}
.spawn(format!("kubernetes-orchestrator-worker:{namespace}"));

Expand Down Expand Up @@ -306,6 +315,7 @@ struct OrchestratorWorker {
command_rx: mpsc::UnboundedReceiver<WorkerCommand>,
name_prefix: String,
collect_pod_metrics: bool,
metrics: OrchestratorMetrics,
}

#[derive(Deserialize, Clone, Debug)]
Expand Down Expand Up @@ -1565,18 +1575,28 @@ impl OrchestratorWorker {

let clusterd_usage_fut = get_clusterd_usage(self_, service_name, i);
let (metrics, clusterd_usage) =
match futures::future::join(self_.metrics_api.get(&name), clusterd_usage_fut).await
{
(Ok(metrics), Ok(clusterd_usage)) => (metrics, Some(clusterd_usage)),
(Ok(metrics), Err(e)) => {
warn!("Failed to fetch clusterd usage for {name}: {e}");
(metrics, None)
}
(Err(e), _) => {
warn!("Failed to get metrics for {name}: {e}");
return ServiceProcessMetrics::default();
}
};
futures::future::join(self_.metrics_api.get(&name), clusterd_usage_fut).await;
let clusterd_usage = match clusterd_usage {
Ok(clusterd_usage) => Some(clusterd_usage),
Err(e) => {
self_
.metrics
.record_fetch_error(FetchStep::ClusterdUsage, e.as_ref());
warn!("Failed to fetch clusterd usage for {name}: {e:#}");
None
}
};
let metrics = match metrics {
Ok(metrics) => metrics,
Err(e) => {
self_.metrics.record_fetch_error(FetchStep::PodMetrics, &e);
warn!(
"Failed to get metrics for {name}: {}",
e.display_with_causes()
);
return ServiceProcessMetrics::default();
}
};
let Some(PodMetricsContainer {
usage:
PodMetricsContainerUsage {
Expand All @@ -1586,6 +1606,9 @@ impl OrchestratorWorker {
..
}) = metrics.containers.get(0)
else {
self_
.metrics
.record_fetch_failure(FetchStep::PodMetrics, "no_containers");
warn!("metrics result contained no containers for {name}");
return ServiceProcessMetrics::default();
};
Expand Down
101 changes: 101 additions & 0 deletions src/orchestrator-kubernetes/src/metrics.rs
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;
52 changes: 52 additions & 0 deletions src/orchestrator-kubernetes/src/metrics/tests.rs
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");
}
13 changes: 12 additions & 1 deletion src/orchestrator-kubernetes/src/util.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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

Copy link
Copy Markdown
Contributor

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".

// 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

@ggevay ggevay Oct 9, 2026 •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The 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 Config::reset_reader_on_write, and it can only ship in a kube-client release after 4.2.0 (we're on 3.1.0). Suggest: "TODO: set reset_reader_on_write and return to 60 s once we're on a kube-client release with kube-rs#2111."

// `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();
Expand Down
Loading