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
17 changes: 12 additions & 5 deletions src/coordinator/query_coordinator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -312,8 +312,8 @@ impl<'a> StageCoordinator<'a> {
stage_id: self.stage_id,
task_number: task_i,
};
let task_metrics = self.metrics_store.clone();
let completed_dynamic_filter_store = self.completed_dynamic_filter_store.clone();
let mut task_metrics = self.metrics_store.clone();
let mut completed_dynamic_filter_store = self.completed_dynamic_filter_store.clone();
let dynamic_filter_registry = Arc::clone(self.dynamic_filter_registry);
let (load_info_tx, load_info_rx) = tokio::sync::mpsc::unbounded_channel();
let mut load_info_tx_opt = Some(load_info_tx);
Expand All @@ -325,8 +325,8 @@ impl<'a> StageCoordinator<'a> {
while let Some(msg) = worker_to_coordinator_rx.recv().await {
match msg {
WorkerToCoordinatorMsg::TaskMetrics(v) => {
if let Some(task_metrics) = &task_metrics {
task_metrics.insert(task_key, v);
if let Some(store) = task_metrics.take() {
store.insert(task_key, v);
}
}
WorkerToCoordinatorMsg::LoadInfo(load_info) => {
Expand All @@ -338,7 +338,7 @@ impl<'a> StageCoordinator<'a> {
let _ = load_info_tx_opt.take();
}
WorkerToCoordinatorMsg::TaskCompletedDynamicFilters(filters) => {
if let Some(store) = &completed_dynamic_filter_store {
if let Some(store) = completed_dynamic_filter_store.take() {
store.insert(task_key, filters);
}
}
Expand All @@ -347,6 +347,13 @@ impl<'a> StageCoordinator<'a> {
}
}
}
// An unexecuted task sends no final reports; still complete its waits.
if let Some(store) = task_metrics {
store.insert(task_key, TaskMetrics::default());
}
if let Some(store) = completed_dynamic_filter_store {
store.insert(task_key, TaskCompletedDynamicFilters::default());
}
});
load_info_rx
}
Expand Down
3 changes: 3 additions & 0 deletions src/metrics/task_metrics_rewriter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -276,6 +276,9 @@ pub fn stage_metrics_rewriter(
stage.num
);
};
if task_metrics.pre_order_plan_metrics.is_empty() {
continue; // The task was never executed, so there are no plan metrics to rewrite.
}
Comment thread
gabotechs marked this conversation as resolved.

let mut per_task_counter = 0usize;
stage.plan.apply_with_dt_ctx(d_ctx, |node, _ctx| {
Expand Down
2 changes: 1 addition & 1 deletion src/protocol/worker_channel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -164,7 +164,7 @@ pub struct TaskDynamicFilter {
pub expression: MaybeEncoded<Arc<dyn PhysicalExpr>>,
}

#[derive(Clone, Debug)]
#[derive(Clone, Debug, Default)]
pub struct TaskMetrics {
/// Metrics for a single task's plan nodes in pre-order traversal order.
/// The TaskKey is implicit — it is determined by the SetPlanRequest that
Expand Down
10 changes: 8 additions & 2 deletions tests/metrics_collection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ mod tests {
use datafusion_distributed::{
DefaultSessionBuilder, DistributedExt, DistributedLeafExec, DistributedMetricsFormat,
NetworkCoalesceExec, NetworkShuffleExec, WorkerQueryContext, display_plan_ascii,
rewrite_distributed_plan_with_metrics,
rewrite_distributed_plan_with_dynamic_filters, rewrite_distributed_plan_with_metrics,
};
use futures::TryStreamExt;
use std::sync::Arc;
Expand Down Expand Up @@ -382,11 +382,11 @@ mod tests {

/// Regression for #739: an empty build side can leave sampled probe tasks unexecuted.
#[tokio::test]
#[ignore = "metrics rewrite hangs on planned but unexecuted tasks"]
async fn metrics_rewrite_after_unexecuted_aqe_tasks() -> Result<(), Box<dyn std::error::Error>>
{
let (mut ctx, _guard, _) = start_localhost_context(3, DefaultSessionBuilder).await;
ctx.set_distributed_dynamic_task_count(true)?;
ctx = ctx.with_distributed_dynamic_filter_collection(true)?;
register_parquet_tables(&ctx).await?;
{
let state = ctx.state_ref();
Expand All @@ -409,6 +409,12 @@ mod tests {
.await?;
assert_eq!(batches.iter().map(|b| b.num_rows()).sum::<usize>(), 0);

let task_ctx = ctx.task_ctx();
let plan = tokio::time::timeout(
Duration::from_secs(10),
rewrite_distributed_plan_with_dynamic_filters(plan, &task_ctx),
)
.await??;
tokio::time::timeout(
Duration::from_secs(10),
rewrite_distributed_plan_with_metrics(plan, DistributedMetricsFormat::PerTask),
Expand Down
Loading