From 6f66aa47dffd4f352ff8736bfbe9abc3298cc754 Mon Sep 17 00:00:00 2001 From: Gabriel Date: Wed, 26 Aug 2026 16:23:00 +0200 Subject: [PATCH] refactor: centralize distributed metric names --- src/coordinator/prepare_dynamic_plan.rs | 14 +-- src/coordinator/query_coordinator.rs | 9 +- src/execution_plans/sampler.rs | 26 ++++-- src/lib.rs | 26 +++++- src/metrics/mod.rs | 96 +++++++++++++++++++++ src/protocol/grpc/metrics_proto.rs | 3 +- src/protocol/grpc/worker_client.rs | 29 ++++--- src/work_unit_feed/remote_work_unit_feed.rs | 24 ++++-- src/worker/task_data.rs | 5 +- src/worker/worker_connection_pool.rs | 3 +- tests/local_connections.rs | 14 +-- tests/metrics_collection.rs | 51 ++++++----- 12 files changed, 225 insertions(+), 75 deletions(-) diff --git a/src/coordinator/prepare_dynamic_plan.rs b/src/coordinator/prepare_dynamic_plan.rs index dff4e285a..43851614b 100644 --- a/src/coordinator/prepare_dynamic_plan.rs +++ b/src/coordinator/prepare_dynamic_plan.rs @@ -7,6 +7,10 @@ use crate::distributed_planner::{ }; use crate::events::TaskCountAnnotation::{Desired, Maximum}; use crate::execution_plans::SamplerExec; +use crate::metrics::{ + CPU_COST_METRIC, ESTIMATED_OUTPUT_BYTES_METRIC, ESTIMATED_PCT_SAMPLED_METRIC, + MEMORY_COST_METRIC, NETWORK_COST_METRIC, +}; use crate::stage::{LocalStage, RemoteStage}; use crate::{ BytesCounterMetric, LoadInfo, MaxGaugeMetric, NetworkBoundaryExt, NetworkCoalesceExec, Stage, @@ -48,15 +52,15 @@ pub(super) async fn prepare_dynamic_plan( // by the SamplerExec injected by this same function. let cost = calculate_cost(&input_stage.plan)?; metrics.push(BytesCounterMetric::new_metric( - "cpu_cost", + CPU_COST_METRIC, *cost.cpu.get_value().unwrap_or(&0), )); metrics.push(BytesCounterMetric::new_metric( - "memory_cost", + MEMORY_COST_METRIC, *cost.memory.get_value().unwrap_or(&0), )); metrics.push(BytesCounterMetric::new_metric( - "network_cost", + NETWORK_COST_METRIC, *cost.network.get_value().unwrap_or(&0), )); let compute_based_task_count = cost @@ -285,7 +289,7 @@ async fn gather_runtime_statistics( }; new_metrics.push(MaxGaugeMetric::new_metric( - "estimated_pct_sampled", + ESTIMATED_PCT_SAMPLED_METRIC, (estimated_pct_sampled * 100.) as usize, )); @@ -299,7 +303,7 @@ async fn gather_runtime_statistics( let total_byte_size: usize = per_col_byte_size.iter().sum(); new_metrics.push(BytesCounterMetric::new_metric( - "estimated_output_bytes", + ESTIMATED_OUTPUT_BYTES_METRIC, total_byte_size, )); diff --git a/src/coordinator/query_coordinator.rs b/src/coordinator/query_coordinator.rs index 2f9c99fb0..0533e1d6a 100644 --- a/src/coordinator/query_coordinator.rs +++ b/src/coordinator/query_coordinator.rs @@ -5,6 +5,9 @@ use crate::coordinator::MetricsStore; use crate::coordinator::latency_metric::LatencyMetric; use crate::events::{RouteTasksEvent, RouteTasksHandlers}; use crate::execution_plans::{ChildrenIsolatorUnionExec, DistributedLeafExec}; +use crate::metrics::{ + LOCAL_COORDINATOR_CHANNELS_METRIC, PLAN_SEND_LATENCY_METRIC, REMOTE_COORDINATOR_CHANNELS_METRIC, +}; use crate::passthrough_headers::get_passthrough_headers; use crate::stage::LocalStage; use crate::work_unit_feed::WorkUnitFeedRegistry; @@ -457,12 +460,12 @@ impl CoordinatorToWorkerMetrics { pub(super) fn new(metrics: &ExecutionPlanMetricsSet) -> Self { Self { local_coordinator_channels: MetricBuilder::new(metrics) - .global_counter("local_coordinator_channels"), + .global_counter(LOCAL_COORDINATOR_CHANNELS_METRIC), remote_coordinator_channels: MetricBuilder::new(metrics) - .global_counter("remote_coordinator_channels"), + .global_counter(REMOTE_COORDINATOR_CHANNELS_METRIC), // Latency statistics about the network calls issued to the workers for feeding subplans. plan_send_latency: Arc::new(LatencyMetric::new( - "plan_send_latency", + PLAN_SEND_LATENCY_METRIC, with_task_id_label, metrics, )), diff --git a/src/execution_plans/sampler.rs b/src/execution_plans/sampler.rs index 71b72437d..7bcd58614 100644 --- a/src/execution_plans/sampler.rs +++ b/src/execution_plans/sampler.rs @@ -1,4 +1,10 @@ use crate::common::{TreeNodeExt, require_one_child, vec_cast}; +use crate::metrics::{ + BYTES_READY_METRIC, KICK_OFF_TO_EXECUTION_MAX_METRIC, KICK_OFF_TO_EXECUTION_P50_METRIC, + KICK_OFF_TO_FIRST_BATCH_MAX_METRIC, KICK_OFF_TO_FIRST_BATCH_P50_METRIC, + KICK_OFF_TO_LOAD_INFO_SENT_MAX_METRIC, KICK_OFF_TO_LOAD_INFO_SENT_P50_METRIC, + MAX_BATCHES_PEEKED_METRIC, MAX_MEMORY_USED_METRIC, +}; use crate::{ BytesCounterMetric, BytesMetricExt, GaugeMetricExt, LatencyMetricExt, LoadInfo, MaxGaugeMetric, MaxLatencyMetric, P50LatencyMetric, @@ -75,15 +81,17 @@ impl SamplerExecMetrics { fn new(metric_set: &ExecutionPlanMetricsSet) -> Self { let bdr = || MetricBuilder::new(metric_set); Self { - kick_off_to_fist_batch_p50: bdr().p50_latency("kick_off_to_first_batch_p50"), - kick_off_to_fist_batch_max: bdr().max_latency("kick_off_to_first_batch_max"), - kick_off_to_load_info_sent_p50: bdr().p50_latency("kick_off_to_load_info_sent_p50"), - kick_off_to_load_info_sent_max: bdr().max_latency("kick_off_to_load_info_sent_max"), - kick_off_to_execution_p50: bdr().p50_latency("kick_off_to_execution_p50"), - kick_off_to_execution_max: bdr().max_latency("kick_off_to_execution_max"), - max_batches_peeked: bdr().max_gauge("max_batches_peeked"), - max_mem_used: bdr().global_gauge("max_mem_used"), - bytes_ready: bdr().bytes_counter("bytes_ready"), + kick_off_to_fist_batch_p50: bdr().p50_latency(KICK_OFF_TO_FIRST_BATCH_P50_METRIC), + kick_off_to_fist_batch_max: bdr().max_latency(KICK_OFF_TO_FIRST_BATCH_MAX_METRIC), + kick_off_to_load_info_sent_p50: bdr() + .p50_latency(KICK_OFF_TO_LOAD_INFO_SENT_P50_METRIC), + kick_off_to_load_info_sent_max: bdr() + .max_latency(KICK_OFF_TO_LOAD_INFO_SENT_MAX_METRIC), + kick_off_to_execution_p50: bdr().p50_latency(KICK_OFF_TO_EXECUTION_P50_METRIC), + kick_off_to_execution_max: bdr().max_latency(KICK_OFF_TO_EXECUTION_MAX_METRIC), + max_batches_peeked: bdr().max_gauge(MAX_BATCHES_PEEKED_METRIC), + max_mem_used: bdr().global_gauge(MAX_MEMORY_USED_METRIC), + bytes_ready: bdr().bytes_counter(BYTES_READY_METRIC), elapsed_compute: { let time = Time::new(); bdr().build(MetricValue::ElapsedCompute(time.clone())); diff --git a/src/lib.rs b/src/lib.rs index 3ad27ddd9..1b19a5db0 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -34,10 +34,28 @@ pub use execution_plans::{ NetworkShuffleExec, }; pub use metrics::{ - AvgLatencyMetric, BytesCounterMetric, BytesMetricExt, DISTRIBUTED_DATAFUSION_TASK_ID_LABEL, - DistributedMetricsFormat, FirstLatencyMetric, GaugeMetricExt, LatencyMetricExt, MaxGaugeMetric, - MaxLatencyMetric, MinLatencyMetric, P50LatencyMetric, P75LatencyMetric, P95LatencyMetric, - P99LatencyMetric, rewrite_distributed_plan_with_metrics, + AvgLatencyMetric, BytesCounterMetric, BytesMetricExt, DistributedMetricsFormat, + FirstLatencyMetric, GaugeMetricExt, LatencyMetricExt, MaxGaugeMetric, MaxLatencyMetric, + MinLatencyMetric, P50LatencyMetric, P75LatencyMetric, P95LatencyMetric, P99LatencyMetric, + rewrite_distributed_plan_with_metrics, +}; +pub use metrics::{ + BYTES_READY_METRIC, BYTES_TRANSFERRED_METRIC, CPU_COST_METRIC, + DISTRIBUTED_DATAFUSION_TASK_ID_LABEL, ESTIMATED_OUTPUT_BYTES_METRIC, + ESTIMATED_PCT_SAMPLED_METRIC, KICK_OFF_TO_EXECUTION_MAX_METRIC, + KICK_OFF_TO_EXECUTION_P50_METRIC, KICK_OFF_TO_FIRST_BATCH_MAX_METRIC, + KICK_OFF_TO_FIRST_BATCH_P50_METRIC, KICK_OFF_TO_LOAD_INFO_SENT_MAX_METRIC, + KICK_OFF_TO_LOAD_INFO_SENT_P50_METRIC, LOCAL_CONNECTIONS_USED_METRIC, + LOCAL_COORDINATOR_CHANNELS_METRIC, MAX_BATCHES_PEEKED_METRIC, MAX_MEMORY_USED_METRIC, + MEMORY_COST_METRIC, MESSAGE_COUNT_METRIC, NETWORK_COST_METRIC, NETWORK_LATENCY_COUNT_METRIC, + NETWORK_LATENCY_FIRST_METRIC, NETWORK_LATENCY_MAX_METRIC, NETWORK_LATENCY_MIN_METRIC, + NETWORK_LATENCY_P50_METRIC, NETWORK_LATENCY_P95_METRIC, NETWORK_LATENCY_SUM_METRIC, + PLAN_ADDED_AT_METRIC, PLAN_BYTES_SENT_METRIC, PLAN_EXECUTED_AT_METRIC, PLAN_FINISHED_AT_METRIC, + PLAN_SEND_LATENCY_METRIC, REMOTE_COORDINATOR_CHANNELS_METRIC, WORK_UNIT_BYTES_METRIC, + WORK_UNIT_COUNT_METRIC, WORK_UNIT_IN_MEMORY_COUNT_METRIC, + WORK_UNIT_PROCESSED_LATENCY_MAX_METRIC, WORK_UNIT_PROCESSED_LATENCY_P50_METRIC, + WORK_UNIT_RECEIVED_LATENCY_MAX_METRIC, WORK_UNIT_RECEIVED_LATENCY_P50_METRIC, + WORK_UNIT_SEND_LATENCY_MAX_METRIC, WORK_UNIT_SEND_LATENCY_P50_METRIC, }; pub use protocol::LocalWorkerContext; diff --git a/src/metrics/mod.rs b/src/metrics/mod.rs index fd1139426..03e846605 100644 --- a/src/metrics/mod.rs +++ b/src/metrics/mod.rs @@ -12,6 +12,102 @@ pub use latency_metric::{ pub use max_gauge_metric::{GaugeMetricExt, MaxGaugeMetric}; pub(crate) use task_metrics_collector::collect_plan_metrics; pub use task_metrics_rewriter::{DistributedMetricsFormat, rewrite_distributed_plan_with_metrics}; + +/// Emitted by dynamic-planner stage records; estimates the CPU cost of the stage input. +pub const CPU_COST_METRIC: &str = "cpu_cost"; +/// Emitted by dynamic-planner stage records; estimates the memory cost of the stage input. +pub const MEMORY_COST_METRIC: &str = "memory_cost"; +/// Emitted by dynamic-planner stage records; estimates the network cost of the stage input. +pub const NETWORK_COST_METRIC: &str = "network_cost"; +/// Emitted by dynamic-planner stage records; estimates the percentage of input sampled. +pub const ESTIMATED_PCT_SAMPLED_METRIC: &str = "estimated_pct_sampled"; +/// Emitted by dynamic-planner stage records; estimates the stage's total output size in bytes. +pub const ESTIMATED_OUTPUT_BYTES_METRIC: &str = "estimated_output_bytes"; +/// Emitted by `DistributedExec`; counts coordinator-to-worker channels routed locally. +pub const LOCAL_COORDINATOR_CHANNELS_METRIC: &str = "local_coordinator_channels"; +/// Emitted by `DistributedExec`; counts coordinator-to-worker channels routed remotely. +pub const REMOTE_COORDINATOR_CHANNELS_METRIC: &str = "remote_coordinator_channels"; +/// Emitted by `DistributedExec`; measures latency for sending a plan to a worker. +pub const PLAN_SEND_LATENCY_METRIC: &str = "plan_send_latency"; +/// Emitted by coordinator-to-worker streams; counts encoded plan bytes sent to workers. +pub const PLAN_BYTES_SENT_METRIC: &str = "plan_bytes_sent"; +/// Emitted by worker task data; records when a coordinator added the task plan. +pub const PLAN_ADDED_AT_METRIC: &str = "plan_added_at"; +/// Emitted by worker task data; records when the worker began executing the task plan. +pub const PLAN_EXECUTED_AT_METRIC: &str = "plan_executed_at"; +/// Emitted by worker task data; records when the worker finished the task plan's stream. +pub const PLAN_FINISHED_AT_METRIC: &str = "plan_finished_at"; + +/// Emitted by `NetworkCoalesceExec`, `NetworkShuffleExec`, and `NetworkBroadcastExec`. +/// Counts encoded record-batch bytes received from workers. +pub const BYTES_TRANSFERRED_METRIC: &str = "bytes_transferred"; +/// Emitted by `NetworkCoalesceExec`, `NetworkShuffleExec`, `NetworkBroadcastExec`, and `SamplerExec`. +/// Records peak buffered memory in bytes. +pub const MAX_MEMORY_USED_METRIC: &str = "max_mem_used"; +/// Emitted by `NetworkCoalesceExec`, `NetworkShuffleExec`, and `NetworkBroadcastExec`. +/// Counts messages received from workers. +pub const MESSAGE_COUNT_METRIC: &str = "msg_count"; +/// Emitted by `NetworkCoalesceExec`, `NetworkShuffleExec`, and `NetworkBroadcastExec`. +/// Records the minimum worker-message latency. +pub const NETWORK_LATENCY_MIN_METRIC: &str = "network_latency_min"; +/// Emitted by `NetworkCoalesceExec`, `NetworkShuffleExec`, and `NetworkBroadcastExec`. +/// Records the maximum worker-message latency. +pub const NETWORK_LATENCY_MAX_METRIC: &str = "network_latency_max"; +/// Emitted by `NetworkCoalesceExec`, `NetworkShuffleExec`, and `NetworkBroadcastExec`. +/// Records the 50th-percentile worker-message latency. +pub const NETWORK_LATENCY_P50_METRIC: &str = "network_latency_p50"; +/// Emitted by `NetworkCoalesceExec`, `NetworkShuffleExec`, and `NetworkBroadcastExec`. +/// Records the 95th-percentile worker-message latency. +pub const NETWORK_LATENCY_P95_METRIC: &str = "network_latency_p95"; +/// Emitted by `NetworkCoalesceExec`, `NetworkShuffleExec`, and `NetworkBroadcastExec`. +/// Records latency of the first worker message. +pub const NETWORK_LATENCY_FIRST_METRIC: &str = "network_latency_first"; +/// Emitted by `NetworkCoalesceExec`, `NetworkShuffleExec`, and `NetworkBroadcastExec`. +/// Sums worker-message latencies. +pub const NETWORK_LATENCY_SUM_METRIC: &str = "network_latency_sum"; +/// Emitted by `NetworkCoalesceExec`, `NetworkShuffleExec`, and `NetworkBroadcastExec`. +/// Counts worker messages included in latency metrics. +pub const NETWORK_LATENCY_COUNT_METRIC: &str = "network_latency_count"; + +/// Emitted by `SamplerExec`; records P50 time from sampler kickoff to its first batch. +pub const KICK_OFF_TO_FIRST_BATCH_P50_METRIC: &str = "kick_off_to_first_batch_p50"; +/// Emitted by `SamplerExec`; records maximum time from sampler kickoff to its first batch. +pub const KICK_OFF_TO_FIRST_BATCH_MAX_METRIC: &str = "kick_off_to_first_batch_max"; +/// Emitted by `SamplerExec`; records P50 time from kickoff to sending load information. +pub const KICK_OFF_TO_LOAD_INFO_SENT_P50_METRIC: &str = "kick_off_to_load_info_sent_p50"; +/// Emitted by `SamplerExec`; records maximum time from kickoff to sending load information. +pub const KICK_OFF_TO_LOAD_INFO_SENT_MAX_METRIC: &str = "kick_off_to_load_info_sent_max"; +/// Emitted by `SamplerExec`; records P50 time from kickoff to execution. +pub const KICK_OFF_TO_EXECUTION_P50_METRIC: &str = "kick_off_to_execution_p50"; +/// Emitted by `SamplerExec`; records maximum time from kickoff to execution. +pub const KICK_OFF_TO_EXECUTION_MAX_METRIC: &str = "kick_off_to_execution_max"; +/// Emitted by `SamplerExec`; records the largest number of record batches held for sampling. +pub const MAX_BATCHES_PEEKED_METRIC: &str = "max_batches_peeked"; +/// Emitted by `SamplerExec`; counts bytes ready when it reports load information. +pub const BYTES_READY_METRIC: &str = "bytes_ready"; + +/// Emitted by `NetworkCoalesceExec`, `NetworkShuffleExec`, and `NetworkBroadcastExec`. +/// Counts worker connections resolved to the local process. +pub const LOCAL_CONNECTIONS_USED_METRIC: &str = "local_connections_used"; +/// Emitted by `RemoteFeedProvider`; counts encoded work-unit bytes received from the coordinator. +pub const WORK_UNIT_BYTES_METRIC: &str = "work_unit_bytes"; +/// Emitted by `RemoteFeedProvider`; counts work units delivered in-memory rather than over transport. +pub const WORK_UNIT_IN_MEMORY_COUNT_METRIC: &str = "work_unit_in_memory_count"; +/// Emitted by `RemoteFeedProvider`; counts work units received from the coordinator. +pub const WORK_UNIT_COUNT_METRIC: &str = "work_unit_count"; +/// Emitted by `RemoteFeedProvider`; records maximum coordinator-to-worker work-unit send latency. +pub const WORK_UNIT_SEND_LATENCY_MAX_METRIC: &str = "work_unit_send_latency_max"; +/// Emitted by `RemoteFeedProvider`; records P50 coordinator-to-worker work-unit send latency. +pub const WORK_UNIT_SEND_LATENCY_P50_METRIC: &str = "work_unit_send_latency_p50"; +/// Emitted by `RemoteFeedProvider`; records maximum work-unit receive latency. +pub const WORK_UNIT_RECEIVED_LATENCY_MAX_METRIC: &str = "work_unit_received_latency_max"; +/// Emitted by `RemoteFeedProvider`; records P50 work-unit receive latency. +pub const WORK_UNIT_RECEIVED_LATENCY_P50_METRIC: &str = "work_unit_received_latency_p50"; +/// Emitted by `RemoteFeedProvider`; records maximum work-unit processing latency. +pub const WORK_UNIT_PROCESSED_LATENCY_MAX_METRIC: &str = "work_unit_processed_latency_max"; +/// Emitted by `RemoteFeedProvider`; records P50 work-unit processing latency. +pub const WORK_UNIT_PROCESSED_LATENCY_P50_METRIC: &str = "work_unit_processed_latency_p50"; + /// Label used to annotate metrics in execution plan nodes with the task in which they were executed. /// Note that the same task id may be used in multiple stages. pub const DISTRIBUTED_DATAFUSION_TASK_ID_LABEL: &str = "task_id"; diff --git a/src/protocol/grpc/metrics_proto.rs b/src/protocol/grpc/metrics_proto.rs index b71ec3e41..0198f9f98 100644 --- a/src/protocol/grpc/metrics_proto.rs +++ b/src/protocol/grpc/metrics_proto.rs @@ -600,6 +600,7 @@ pub fn metric_proto_to_df(metric: pb::Metric) -> Result, DataFusionE #[cfg(test)] mod tests { use super::*; + use crate::metrics::BYTES_TRANSFERRED_METRIC; use datafusion::physical_plan::metrics::CustomMetricValue; use datafusion::physical_plan::metrics::{Count, Gauge, Label, MetricsSet, Time, Timestamp}; use datafusion::physical_plan::metrics::{Metric, MetricValue}; @@ -1320,7 +1321,7 @@ mod tests { let mut metrics_set = MetricsSet::new(); metrics_set.push(Arc::new(Metric::new( MetricValue::Custom { - name: Cow::Borrowed("bytes_transferred"), + name: Cow::Borrowed(BYTES_TRANSFERRED_METRIC), value: Arc::new(BytesCounterMetric::from_value(1_073_741_824)), }, Some(0), diff --git a/src/protocol/grpc/worker_client.rs b/src/protocol/grpc/worker_client.rs index 060d8e875..1ea76ff46 100644 --- a/src/protocol/grpc/worker_client.rs +++ b/src/protocol/grpc/worker_client.rs @@ -5,6 +5,12 @@ use super::metrics_proto::metrics_set_proto_to_df; use crate::common::serialize_uuid; use crate::grpc::generated::worker::FlightAppMetadata; use crate::grpc::on_drop_stream::on_drop_stream; +use crate::metrics::{ + BYTES_TRANSFERRED_METRIC, MAX_MEMORY_USED_METRIC, MESSAGE_COUNT_METRIC, + NETWORK_LATENCY_COUNT_METRIC, NETWORK_LATENCY_FIRST_METRIC, NETWORK_LATENCY_MAX_METRIC, + NETWORK_LATENCY_MIN_METRIC, NETWORK_LATENCY_P50_METRIC, NETWORK_LATENCY_P95_METRIC, + NETWORK_LATENCY_SUM_METRIC, PLAN_BYTES_SENT_METRIC, +}; use crate::{ BytesMetricExt, CoordinatorToWorkerMsg, DISTRIBUTED_DATAFUSION_TASK_ID_LABEL, DistributedConfig, ExecuteTaskRequest, FirstLatencyMetric, GetWorkerInfoRequest, @@ -80,7 +86,7 @@ impl WorkerChannel for pb::worker_service_client::WorkerServiceClient Self { - let min_latency = MetricBuilder::new(metrics).min_latency("network_latency_min"); - let max_latency = MetricBuilder::new(metrics).max_latency("network_latency_max"); - let p50_latency = MetricBuilder::new(metrics).p50_latency("network_latency_p50"); - let p95_latency = MetricBuilder::new(metrics).p95_latency("network_latency_p95"); - let first_latency = MetricBuilder::new(metrics).first_latency("network_latency_first"); + let min_latency = MetricBuilder::new(metrics).min_latency(NETWORK_LATENCY_MIN_METRIC); + let max_latency = MetricBuilder::new(metrics).max_latency(NETWORK_LATENCY_MAX_METRIC); + let p50_latency = MetricBuilder::new(metrics).p50_latency(NETWORK_LATENCY_P50_METRIC); + let p95_latency = MetricBuilder::new(metrics).p95_latency(NETWORK_LATENCY_P95_METRIC); + let first_latency = MetricBuilder::new(metrics).first_latency(NETWORK_LATENCY_FIRST_METRIC); let sum_latency = Time::new(); MetricBuilder::new(metrics).build(MetricValue::Time { - name: Cow::Borrowed("network_latency_sum"), + name: Cow::Borrowed(NETWORK_LATENCY_SUM_METRIC), time: sum_latency.clone(), }); - let latency_count = MetricBuilder::new(metrics).counter("network_latency_count", 0); + let latency_count = MetricBuilder::new(metrics).counter(NETWORK_LATENCY_COUNT_METRIC, 0); Self { min_latency, diff --git a/src/work_unit_feed/remote_work_unit_feed.rs b/src/work_unit_feed/remote_work_unit_feed.rs index 5b73ef738..f533b7e8a 100644 --- a/src/work_unit_feed/remote_work_unit_feed.rs +++ b/src/work_unit_feed/remote_work_unit_feed.rs @@ -1,4 +1,10 @@ use crate::common::now_ns; +use crate::metrics::{ + WORK_UNIT_BYTES_METRIC, WORK_UNIT_COUNT_METRIC, WORK_UNIT_IN_MEMORY_COUNT_METRIC, + WORK_UNIT_PROCESSED_LATENCY_MAX_METRIC, WORK_UNIT_PROCESSED_LATENCY_P50_METRIC, + WORK_UNIT_RECEIVED_LATENCY_MAX_METRIC, WORK_UNIT_RECEIVED_LATENCY_P50_METRIC, + WORK_UNIT_SEND_LATENCY_MAX_METRIC, WORK_UNIT_SEND_LATENCY_P50_METRIC, +}; use crate::{ BytesMetricExt, CoordinatorToWorkerMsg, LatencyMetricExt, MaybeEncoded, WorkUnit, WorkUnitBatch, WorkUnitMsg, @@ -123,18 +129,18 @@ impl RemoteFeedProvider { ) -> Result>> { let bdr = || MetricBuilder::new(&self.metrics); - let bytes_transferred = bdr().bytes_counter("work_unit_bytes"); - let in_memory_transferred = bdr().global_counter("work_unit_in_memory_count"); - let msg_count = bdr().global_counter("work_unit_count"); + let bytes_transferred = bdr().bytes_counter(WORK_UNIT_BYTES_METRIC); + let in_memory_transferred = bdr().global_counter(WORK_UNIT_IN_MEMORY_COUNT_METRIC); + let msg_count = bdr().global_counter(WORK_UNIT_COUNT_METRIC); // Track end-to-end network latency distribution for all work units. - let send_latency_max = bdr().max_latency("work_unit_send_latency_max"); - let send_latency_p50 = bdr().p50_latency("work_unit_send_latency_p50"); + let send_latency_max = bdr().max_latency(WORK_UNIT_SEND_LATENCY_MAX_METRIC); + let send_latency_p50 = bdr().p50_latency(WORK_UNIT_SEND_LATENCY_P50_METRIC); - let received_latency_max = bdr().max_latency("work_unit_received_latency_max"); - let received_latency_p50 = bdr().p50_latency("work_unit_received_latency_p50"); + let received_latency_max = bdr().max_latency(WORK_UNIT_RECEIVED_LATENCY_MAX_METRIC); + let received_latency_p50 = bdr().p50_latency(WORK_UNIT_RECEIVED_LATENCY_P50_METRIC); - let processed_latency_max = bdr().max_latency("work_unit_processed_latency_max"); - let processed_latency_p50 = bdr().p50_latency("work_unit_processed_latency_p50"); + let processed_latency_max = bdr().max_latency(WORK_UNIT_PROCESSED_LATENCY_MAX_METRIC); + let processed_latency_p50 = bdr().p50_latency(WORK_UNIT_PROCESSED_LATENCY_P50_METRIC); let elapsed_compute = bdr().elapsed_compute(partition); diff --git a/src/worker/task_data.rs b/src/worker/task_data.rs index 2ff377173..cba76aff8 100644 --- a/src/worker/task_data.rs +++ b/src/worker/task_data.rs @@ -1,5 +1,6 @@ use crate::common::OnceLockResult; use crate::common::now_ns; +use crate::metrics::{PLAN_ADDED_AT_METRIC, PLAN_EXECUTED_AT_METRIC, PLAN_FINISHED_AT_METRIC}; use crate::{MaxLatencyMetric, ProducerHead, TaskMetrics}; use datafusion::common::{DataFusionError, Result}; use datafusion::execution::TaskContext; @@ -28,10 +29,6 @@ pub struct TaskData { pub(super) task_data_metrics: Arc, } -pub(crate) const PLAN_ADDED_AT_METRIC: &str = "plan_added_at"; -pub(crate) const PLAN_EXECUTED_AT_METRIC: &str = "plan_executed_at"; -pub(crate) const PLAN_FINISHED_AT_METRIC: &str = "plan_finished_at"; - #[derive(Debug)] pub(super) struct TaskDataMetrics { pub(super) query_start_time_ns: usize, diff --git a/src/worker/worker_connection_pool.rs b/src/worker/worker_connection_pool.rs index 6f540655c..f83ad1419 100644 --- a/src/worker/worker_connection_pool.rs +++ b/src/worker/worker_connection_pool.rs @@ -1,4 +1,5 @@ use crate::distributed_planner::ProducerHead; +use crate::metrics::LOCAL_CONNECTIONS_USED_METRIC; use crate::passthrough_headers::get_passthrough_headers; use crate::stage::RemoteStage; use crate::{ExecuteTaskRequest, LocalWorkerContext, TaskKey, get_distributed_channel_resolver}; @@ -63,7 +64,7 @@ impl WorkerConnectionPool { let bdr = || MetricBuilder::new(&self.metrics); let output_bytes = bdr().output_bytes(target_partition); let output_rows = bdr().output_rows(target_partition); - let local_connections_used = bdr().global_counter("local_connections_used"); + let local_connections_used = bdr().global_counter(LOCAL_CONNECTIONS_USED_METRIC); // Otherwise, we need to reach the remote worker through the `WorkerChannel`. Unlike local // connections, these remote connections span a range of partitions so that `WorkerChannel` diff --git a/tests/local_connections.rs b/tests/local_connections.rs index 84648b3cd..ac61aa2f1 100644 --- a/tests/local_connections.rs +++ b/tests/local_connections.rs @@ -6,8 +6,10 @@ mod tests { use datafusion_distributed::test_utils::localhost::start_localhost_context; use datafusion_distributed::test_utils::parquet::register_parquet_tables; use datafusion_distributed::{ - DefaultSessionBuilder, DistributedExt, DistributedMetricsFormat, NetworkBoundaryExt, - display_plan_ascii, rewrite_distributed_plan_with_metrics, + DefaultSessionBuilder, DistributedExt, DistributedMetricsFormat, + LOCAL_CONNECTIONS_USED_METRIC, LOCAL_COORDINATOR_CHANNELS_METRIC, NetworkBoundaryExt, + REMOTE_COORDINATOR_CHANNELS_METRIC, display_plan_ascii, + rewrite_distributed_plan_with_metrics, }; use std::sync::Arc; @@ -37,7 +39,7 @@ mod tests { let metrics = plan.metrics().unwrap(); let local_connections_used = metrics - .sum(|v| v.value().name() == "local_connections_used") + .sum(|v| v.value().name() == LOCAL_CONNECTIONS_USED_METRIC) .map_or(0, |v| v.as_usize()); if local_connections_used == 0 { return internal_err!("local_connections_used==0"); @@ -74,8 +76,8 @@ mod tests { // - coordinator worker 0 -> stage 0, worker 1 | remote // - coordinator worker 0 -> stage 1, worker 0 | local // - coordinator worker 0 -> stage 1, worker 1 | remote - assert_eq!(metric_value("local_coordinator_channels"), 2); - assert_eq!(metric_value("remote_coordinator_channels"), 2); + assert_eq!(metric_value(LOCAL_COORDINATOR_CHANNELS_METRIC), 2); + assert_eq!(metric_value(REMOTE_COORDINATOR_CHANNELS_METRIC), 2); Ok(()) } @@ -100,7 +102,7 @@ mod tests { local_connections_used += node .metrics() .unwrap_or_default() - .sum(|metric| metric.value().name() == "local_connections_used") + .sum(|metric| metric.value().name() == LOCAL_CONNECTIONS_USED_METRIC) .map_or(0, |metric| metric.as_usize()); } Ok(TreeNodeRecursion::Continue) diff --git a/tests/metrics_collection.rs b/tests/metrics_collection.rs index c6a421da7..fc2d64a17 100644 --- a/tests/metrics_collection.rs +++ b/tests/metrics_collection.rs @@ -16,9 +16,12 @@ mod tests { row_generator_desired_task_count_handler, row_generator_scale_up_leaf_node_handler, }; use datafusion_distributed::{ - DefaultSessionBuilder, DistributedExt, DistributedLeafExec, DistributedMetricsFormat, - NetworkCoalesceExec, NetworkShuffleExec, WorkerQueryContext, display_plan_ascii, - rewrite_distributed_plan_with_metrics, + BYTES_TRANSFERRED_METRIC, DefaultSessionBuilder, DistributedExt, DistributedLeafExec, + DistributedMetricsFormat, MAX_MEMORY_USED_METRIC, NETWORK_LATENCY_COUNT_METRIC, + NETWORK_LATENCY_FIRST_METRIC, NETWORK_LATENCY_MAX_METRIC, NETWORK_LATENCY_MIN_METRIC, + NETWORK_LATENCY_P50_METRIC, NETWORK_LATENCY_SUM_METRIC, NetworkCoalesceExec, + NetworkShuffleExec, PLAN_ADDED_AT_METRIC, PLAN_EXECUTED_AT_METRIC, PLAN_FINISHED_AT_METRIC, + WorkerQueryContext, display_plan_ascii, rewrite_distributed_plan_with_metrics, }; use futures::TryStreamExt; use std::sync::Arc; @@ -161,42 +164,46 @@ mod tests { println!("{}", display_plan_ascii(s_physical.as_ref(), true)); println!("{}", display_plan_ascii(d_physical.as_ref(), true)); - let value = node_metrics::(&d_physical, "bytes_transferred", 1); + let value = node_metrics::(&d_physical, BYTES_TRANSFERRED_METRIC, 1); assert!(value > 100); - let value = node_metrics::(&d_physical, "max_mem_used", 1); + let value = node_metrics::(&d_physical, MAX_MEMORY_USED_METRIC, 1); assert!(value > 100); let value = node_metrics::(&d_physical, "elapsed_compute", 1); assert!(value > 100); - let value = node_metrics::(&d_physical, "network_latency_min", 1); + let value = node_metrics::(&d_physical, NETWORK_LATENCY_MIN_METRIC, 1); assert!(value > 0); - let value = node_metrics::(&d_physical, "network_latency_max", 1); + let value = node_metrics::(&d_physical, NETWORK_LATENCY_MAX_METRIC, 1); assert!(value > 0); - let value = node_metrics::(&d_physical, "network_latency_p50", 1); + let value = node_metrics::(&d_physical, NETWORK_LATENCY_P50_METRIC, 1); assert!(value > 0); - let value = node_metrics::(&d_physical, "network_latency_first", 1); + let value = + node_metrics::(&d_physical, NETWORK_LATENCY_FIRST_METRIC, 1); assert!(value > 0); - let value = node_metrics::(&d_physical, "network_latency_sum", 1); + let value = node_metrics::(&d_physical, NETWORK_LATENCY_SUM_METRIC, 1); assert!(value > 0); - let value = node_metrics::(&d_physical, "network_latency_count", 1); + let value = + node_metrics::(&d_physical, NETWORK_LATENCY_COUNT_METRIC, 1); assert!(value > 0); - let value = node_metrics::(&d_physical, "bytes_transferred", 1); + let value = node_metrics::(&d_physical, BYTES_TRANSFERRED_METRIC, 1); assert!(value > 100); - let value = node_metrics::(&d_physical, "max_mem_used", 1); + let value = node_metrics::(&d_physical, MAX_MEMORY_USED_METRIC, 1); assert!(value > 100); let value = node_metrics::(&d_physical, "elapsed_compute", 1); assert!(value > 100); - let value = node_metrics::(&d_physical, "network_latency_min", 1); + let value = node_metrics::(&d_physical, NETWORK_LATENCY_MIN_METRIC, 1); assert!(value > 0); - let value = node_metrics::(&d_physical, "network_latency_max", 1); + let value = node_metrics::(&d_physical, NETWORK_LATENCY_MAX_METRIC, 1); assert!(value > 0); - let value = node_metrics::(&d_physical, "network_latency_p50", 1); + let value = node_metrics::(&d_physical, NETWORK_LATENCY_P50_METRIC, 1); assert!(value > 0); - let value = node_metrics::(&d_physical, "network_latency_first", 1); + let value = + node_metrics::(&d_physical, NETWORK_LATENCY_FIRST_METRIC, 1); assert!(value > 0); - let value = node_metrics::(&d_physical, "network_latency_sum", 1); + let value = node_metrics::(&d_physical, NETWORK_LATENCY_SUM_METRIC, 1); assert!(value > 0); - let value = node_metrics::(&d_physical, "network_latency_count", 1); + let value = + node_metrics::(&d_physical, NETWORK_LATENCY_COUNT_METRIC, 1); assert!(value > 0); Ok(()) @@ -216,9 +223,9 @@ mod tests { let display = display_plan_ascii(d_physical.as_ref(), true); assert_not_contains!(&display, "metrics=[]"); - assert_contains!(&display, "plan_added_at"); - assert_contains!(&display, "plan_executed_at"); - assert_contains!(&display, "plan_finished_at"); + assert_contains!(&display, PLAN_ADDED_AT_METRIC); + assert_contains!(&display, PLAN_EXECUTED_AT_METRIC); + assert_contains!(&display, PLAN_FINISHED_AT_METRIC); Ok(()) }