diff --git a/lib/collection/src/shards/local_shard/telemetry.rs b/lib/collection/src/shards/local_shard/telemetry.rs index ecf8114f58..242da52fb8 100644 --- a/lib/collection/src/shards/local_shard/telemetry.rs +++ b/lib/collection/src/shards/local_shard/telemetry.rs @@ -70,7 +70,7 @@ impl LocalShard { optimizations: OptimizerTelemetry { status, optimizations, - log: (detail.level >= DetailsLevel::Level4) + log: (detail.level >= DetailsLevel::Level4 || detail.optimizer_logs) .then(|| self.optimizers_log.lock().to_telemetry()), }, async_scorer: Some(get_async_scorer()), diff --git a/lib/collection/src/telemetry.rs b/lib/collection/src/telemetry.rs index d3c04a4871..45828d9cd0 100644 --- a/lib/collection/src/telemetry.rs +++ b/lib/collection/src/telemetry.rs @@ -6,6 +6,7 @@ use segment::types::{HnswConfig, Payload, QuantizationConfig, StrictModeConfigOu use serde::Serialize; use uuid::Uuid; +use crate::collection_manager::optimizers::TrackerStatus; use crate::config::{CollectionConfigInternal, CollectionParams, WalConfig}; use crate::operations::types::{OptimizersStatus, ReshardingInfo, ShardTransferInfo}; use crate::optimizers_builder::OptimizersConfig; @@ -51,6 +52,20 @@ impl CollectionTelemetry { .map(|x| x.num_vectors.unwrap_or(0)) .sum() } + + /// Amount of optimizers currently running. + /// + /// Note: A `DetailsLevel` of 4 or setting `telemetry_detail.optimizer_logs` to true is required. + /// Otherwise, this function will return 0, which may not be correct. + pub fn count_optimizers_running(&self) -> usize { + self.shards + .iter() + .flatten() + .filter_map(|replica_set| replica_set.local.as_ref()) + .flat_map(|local_shard| local_shard.optimizations.log.iter().flatten()) + .filter(|log| log.status == TrackerStatus::Optimizing) + .count() + } } #[derive(Serialize, Clone, Debug, JsonSchema)] diff --git a/lib/common/common/src/types.rs b/lib/common/common/src/types.rs index 6d05026b59..c664aca258 100644 --- a/lib/common/common/src/types.rs +++ b/lib/common/common/src/types.rs @@ -33,6 +33,7 @@ impl PartialOrd for ScoredPointOffset { pub struct TelemetryDetail { pub level: DetailsLevel, pub histograms: bool, + pub optimizer_logs: bool, } #[derive(Copy, Clone, Debug, PartialEq, Eq, PartialOrd, Ord)] @@ -69,6 +70,7 @@ impl Default for TelemetryDetail { TelemetryDetail { level: DetailsLevel::Level0, histograms: false, + optimizer_logs: false, } } } diff --git a/src/actix/api/service_api.rs b/src/actix/api/service_api.rs index 6a73e95ad6..be0791d4ee 100644 --- a/src/actix/api/service_api.rs +++ b/src/actix/api/service_api.rs @@ -46,6 +46,7 @@ fn telemetry( let detail = TelemetryDetail { level: details_level, histograms: false, + optimizer_logs: false, }; let telemetry_collector = telemetry_collector.lock().await; let telemetry_data = telemetry_collector.prepare_data(&access, detail).await; @@ -81,6 +82,7 @@ async fn metrics( TelemetryDetail { level: DetailsLevel::Level3, histograms: true, + optimizer_logs: true, }, ) .await; diff --git a/src/common/metrics.rs b/src/common/metrics.rs index 167e1d22d2..4cae669a06 100644 --- a/src/common/metrics.rs +++ b/src/common/metrics.rs @@ -174,6 +174,26 @@ impl MetricsProvider for CollectionsTelemetry { MetricType::GAUGE, vec![gauge(vector_count as f64, &[])], )); + + let mut total_optimizations_running = 0; + + for collection in self.collections.iter().flatten() { + let collection = match collection { + CollectionTelemetryEnum::Full(collection_telemetry) => collection_telemetry, + CollectionTelemetryEnum::Aggregated(_) => { + continue; + } + }; + + total_optimizations_running += collection.count_optimizers_running(); + } + + metrics.push(metric_family( + "optimizer_running_processes", + "number of currently running optimization processes", + MetricType::GAUGE, + vec![gauge(total_optimizations_running as f64, &[])], + )); } } diff --git a/src/common/telemetry_reporting.rs b/src/common/telemetry_reporting.rs index d7663b20dd..3d353f2959 100644 --- a/src/common/telemetry_reporting.rs +++ b/src/common/telemetry_reporting.rs @@ -12,6 +12,7 @@ use super::telemetry::TelemetryCollector; const DETAIL: TelemetryDetail = TelemetryDetail { level: DetailsLevel::Level2, histograms: false, + optimizer_logs: false, }; const REPORTING_INTERVAL: Duration = Duration::from_secs(60 * 60); // One hour