Skip to content
Closed
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
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ tokio = { version = "1.48", features = ["full"] }
http = "1.3.1"
itertools = "0.14.0"
futures = "0.3.31"
log = "0.4.27"
url = "2.5.7"
uuid = "1.19"
delegate = "0.13.4"
Expand Down
29 changes: 22 additions & 7 deletions docs/source/user-guide/05-metrics.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,21 +8,35 @@ display them.
Distributed DataFusion does this for you, and exposes two functions you can use to build your own
EXPLAIN ANALYZE in application code.

## Enabling collection
## Configuring collection

Metrics collection across network boundaries is **on by default**. You can toggle it explicitly:
Metrics collection across network boundaries is **on by default**. To disable it, call
`.with_distributed_metrics_collection(false)?` on the session-state builder.

When enabled, each worker streams the metrics for its tasks back to the coordinator on a dedicated
channel, so they are not lost even if the result stream is dropped early (for example by a `LIMIT`).

Post-query cleanup, metrics collection, and completed dynamic-filter collection share one deadline
starting when the result stream ends. The default is five seconds. To adjust it, configure
`DistributedConfig::metrics_finalization_timeout_ms` before planning the query:

```rust
use datafusion::execution::SessionStateBuilder;
use datafusion_distributed::{DistributedConfig, DistributedExt, SessionStateBuilderExt};

let state = SessionStateBuilder::new()
.with_default_features()
.with_distributed_option_extension(DistributedConfig {
metrics_finalization_timeout_ms: 10_000,
..Default::default()
})
.with_distributed_worker_resolver(/* ... */)
.with_distributed_planner()
.with_distributed_metrics_collection(true) // default is true
.build();
```

When enabled, each worker streams the metrics for its tasks back to the coordinator on a dedicated
channel, so they are not lost even if the result stream is dropped early (for example by a `LIMIT`).
Set it to zero to avoid waiting for late reports. An incomplete metrics result still returns an
error rather than treating missing work as zero.

## Rendering a plan with metrics

Expand All @@ -33,8 +47,9 @@ These functions, all exported from the crate root, do the work:
used to execute the plan. When displaying both dynamic filters and metrics, apply the
dynamic-filter rewrite first.
- `rewrite_distributed_plan_with_metrics(plan, format)` — folds every task's metrics back into the
coordinator's copy of the plan. It waits for all worker metrics to arrive, so the result is always
complete. The `format` is a `DistributedMetricsFormat`:
coordinator's copy of the plan. It waits for complete worker reports and returns an error if a
channel closes without a report or finalization exceeds the configured deadline. The `format` is a
`DistributedMetricsFormat`:
- `Aggregated` — metrics from all tasks of a stage are summed/aggregated into one value per node.
- `PerTask` — each metric collects its per-task values into a map keyed by task id
(`output_rows={0:.., 1:..}`) so you can see each task individually.
Expand Down
164 changes: 155 additions & 9 deletions src/coordinator/distributed.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,13 +3,14 @@ use crate::coordinator::prepare_dynamic_plan::prepare_dynamic_plan;
use crate::coordinator::prepare_static_plan::prepare_static_plan;
use crate::coordinator::query_coordinator::QueryCoordinator;
use crate::coordinator::store::{Store, task_keys_for_plan};
use crate::distributed_planner::DEFAULT_METRICS_FINALIZATION_TIMEOUT_MS;
use crate::dynamic_filtering::{
is_dynamic_filtering_enabled, sever_dynamic_filter_relationships_in_plan_for_display,
};
use crate::{DistributedConfig, TaskCompletedDynamicFilters, TaskKey, TaskMetrics};
use datafusion::common::internal_datafusion_err;
use datafusion::common::tree_node::TreeNodeRecursion;
use datafusion::common::{HashMap, Result, exec_err};
use datafusion::common::{HashMap, Result, exec_datafusion_err, exec_err};
use datafusion::execution::{SendableRecordBatchStream, TaskContext};
use datafusion::physical_expr::PhysicalExpr;
use datafusion::physical_expr_common::metrics::MetricsSet;
Expand All @@ -19,6 +20,9 @@ use datafusion::physical_plan::{DisplayAs, DisplayFormatType, ExecutionPlan, Pla
use futures::StreamExt;
use std::fmt::Formatter;
use std::sync::{Arc, OnceLock};
use std::time::Duration;
use tokio::sync::watch;
use tokio::time::Instant;

/// [ExecutionPlan] that executes the inner plan in distributed mode.
/// Before executing it, two modifications are lazily performed on the plan:
Expand Down Expand Up @@ -50,6 +54,9 @@ pub struct DistributedExec {
pub(crate) metrics_store: Option<Arc<Store<TaskMetrics>>>,
/// Storage for the completed dynamic filters reported by each worker task.
pub(crate) completed_dynamic_filter_store: Option<Arc<Store<TaskCompletedDynamicFilters>>>,
/// Set when the result stream ends; all post-query waits use the same deadline.
finalization_deadline: watch::Sender<Option<Instant>>,
finalization_timeout: Duration,
}

#[derive(Debug, Clone)]
Expand All @@ -62,12 +69,15 @@ pub(super) struct PreparedPlan {

impl DistributedExec {
pub fn new(base_plan: Arc<dyn ExecutionPlan>) -> Self {
let (finalization_deadline, _) = watch::channel(None);
Self {
base_plan,
prepared_plan: Arc::new(OnceLock::new()),
metrics: ExecutionPlanMetricsSet::new(),
metrics_store: None,
completed_dynamic_filter_store: None,
finalization_deadline,
finalization_timeout: Duration::from_millis(DEFAULT_METRICS_FINALIZATION_TIMEOUT_MS),
}
}

Expand All @@ -89,12 +99,29 @@ impl DistributedExec {
self
}

/// Waits until all worker tasks have reported their metrics back via the coordinator channel
/// if metrics collection is enabled.
pub(crate) fn with_finalization_timeout(mut self, timeout: Duration) -> Self {
self.finalization_timeout = timeout;
self
}

async fn finalization_deadline(&self) -> Instant {
wait_for_finalization_deadline(self.finalization_deadline.subscribe()).await
}

/// Waits for complete metrics, if collection is enabled and execution has been prepared.
/// Returns `None` if a task failed to report or the finalization timeout elapsed.
pub async fn wait_for_metrics(&self) -> Option<HashMap<TaskKey, TaskMetrics>> {
let task_metrics = self.metrics_store.as_ref()?;
let plan = &self.prepared_plan.get()?.plan_for_viz;
Some(task_metrics.wait_for(&task_keys_for_plan(plan)).await)
self.complete_metrics().await.ok()
}

pub(crate) async fn complete_metrics(&self) -> Result<HashMap<TaskKey, TaskMetrics>> {
let store = self
.metrics_store
.as_ref()
.ok_or_else(|| exec_datafusion_err!("metrics collection is disabled"))?;
let plan = &self.prepared_plan()?.plan_for_viz;
let keys = task_keys_for_plan(plan);
wait_for_complete_metrics(store, &keys, self.finalization_deadline().await).await
}

/// Waits until all worker tasks have reported their completed dynamic filters back via
Expand All @@ -104,7 +131,12 @@ impl DistributedExec {
) -> Option<HashMap<TaskKey, TaskCompletedDynamicFilters>> {
let store = self.completed_dynamic_filter_store.as_ref()?;
let plan = &self.prepared_plan.get()?.plan_for_viz;
Some(store.wait_for(&task_keys_for_plan(plan)).await)
tokio::time::timeout_at(
self.finalization_deadline().await,
store.wait_for(&task_keys_for_plan(plan)),
)
.await
.ok()
}

fn prepared_plan(&self) -> Result<PreparedPlan> {
Expand Down Expand Up @@ -152,10 +184,41 @@ impl DistributedExec {
metrics: self.metrics.clone(),
metrics_store: self.metrics_store.clone(),
completed_dynamic_filter_store: self.completed_dynamic_filter_store.clone(),
finalization_deadline: self.finalization_deadline.clone(),
finalization_timeout: self.finalization_timeout,
}))
}
}

async fn wait_for_finalization_deadline(mut rx: watch::Receiver<Option<Instant>>) -> Instant {
(*rx.wait_for(Option::is_some)
.await
.expect("DistributedExec owns the finalization deadline sender"))
.expect("deadline was set before the receiver woke")
}

async fn wait_for_complete_metrics(
store: &Store<TaskMetrics>,
keys: &[TaskKey],
deadline: Instant,
) -> Result<HashMap<TaskKey, TaskMetrics>> {
let snapshot = tokio::time::timeout_at(deadline, store.wait_for_terminal(keys))
.await
.map_err(|_| {
exec_datafusion_err!(
"timed out waiting for worker task metrics; pending: {:?}",
store.snapshot(keys).pending
)
})?;
if !snapshot.terminal_without_report.is_empty() {
return exec_err!(
"worker task metrics missing after coordinator channel closed: {:?}",
snapshot.terminal_without_report
);
}
Ok(snapshot.reported)
}

impl DisplayAs for DistributedExec {
fn fmt_as(&self, _: DisplayFormatType, f: &mut Formatter) -> std::fmt::Result {
write!(f, "DistributedExec")
Expand Down Expand Up @@ -200,6 +263,8 @@ impl ExecutionPlan for DistributedExec {
metrics: self.metrics.clone(),
metrics_store: self.metrics_store.clone(),
completed_dynamic_filter_store: self.completed_dynamic_filter_store.clone(),
finalization_deadline: self.finalization_deadline.clone(),
finalization_timeout: self.finalization_timeout,
}))
}

Expand All @@ -221,6 +286,8 @@ impl ExecutionPlan for DistributedExec {
let base_plan = Arc::clone(&self.base_plan);
let prepared_plan = Arc::clone(&self.prepared_plan);
let collect_dynamic_filters = self.completed_dynamic_filter_store.is_some();
let finalization_deadline = self.finalization_deadline.clone();
let finalization_timeout = self.finalization_timeout;

let query_coordinator = Arc::new(QueryCoordinator::new(
Arc::clone(&context),
Expand All @@ -243,7 +310,8 @@ impl ExecutionPlan for DistributedExec {
// 4. The coordinator->worker channel EOS is received in `impl_coordinator_channel.rs`.
// 5. The metrics are send back in the worker->coordinator channel, and then that
// channel is closed.
let guard = query_coordinator.end_query_guard();
let guard = query_coordinator
.end_query_guard(finalization_deadline.clone(), finalization_timeout);

let d_cfg = DistributedConfig::from_config_options(context.session_config().options())?;
let mut prepared = match d_cfg.dynamic_task_count {
Expand Down Expand Up @@ -271,7 +339,16 @@ impl ExecutionPlan for DistributedExec {
}
drop(guard);
drop(tx);
query_coordinator.drain_pending_tasks().await?;
// DataFusion does not close the result stream until this spawned task returns.
// A half-open coordinator channel can otherwise hang result collection itself,
// even if the later metrics rewrite has its own timeout. Cancel stalled
// background tasks on timeout; their metric receivers then become terminal
// without a report, so accounting remains explicitly incomplete.
query_coordinator
.drain_pending_tasks(
wait_for_finalization_deadline(finalization_deadline.subscribe()).await,
)
.await?;
Ok(())
});

Expand All @@ -282,3 +359,72 @@ impl ExecutionPlan for DistributedExec {
Some(self.metrics.clone_inner())
}
}

#[cfg(test)]
mod tests {
use super::*;
use datafusion::physical_plan::metrics::MetricsSet;
use uuid::Uuid;

#[tokio::test]
async fn a_never_closing_channel_cannot_block_complete_metrics_indefinitely() {
let store = Store::<TaskMetrics>::new();
let key = TaskKey {
query_id: Uuid::new_v4(),
stage_id: 1,
task_number: 0,
};
let deadline = Instant::now() + Duration::from_millis(25);
let result = wait_for_complete_metrics(&store, &[key], deadline).await;
assert!(
result
.unwrap_err()
.to_string()
.contains("timed out waiting")
);
assert_eq!(store.snapshot(&[key]).pending, vec![key]);
// Rewriting after the drain used up the budget must not start another wait.
let result = tokio::time::timeout(
Duration::from_millis(100),
wait_for_complete_metrics(&store, &[key], deadline),
)
.await
.expect("expired deadline should return without another timeout window");
assert!(result.is_err());
}

#[tokio::test]
async fn closed_channel_without_metrics_is_an_error_not_zero_metrics() {
let store = Store::<TaskMetrics>::new();
let query_id = Uuid::new_v4();
let reported = TaskKey {
query_id,
stage_id: 1,
task_number: 0,
};
let lost = TaskKey {
query_id,
stage_id: 1,
task_number: 2,
};
store.insert(
reported,
TaskMetrics {
pre_order_plan_metrics: vec![],
task_metrics: MetricsSet::new(),
},
);
store.mark_terminal(lost);
let result = wait_for_complete_metrics(
&store,
&[reported, lost],
Instant::now() + Duration::from_secs(1),
)
.await;
assert!(result.unwrap_err().to_string().contains("metrics missing"));
let snapshot = store.snapshot(&[reported, lost]);
assert!(snapshot.pending.is_empty());
assert_eq!(snapshot.terminal_without_report, vec![lost]);
assert_eq!(snapshot.reported.len(), 1);
}
}
Loading