Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
14 changes: 14 additions & 0 deletions src/coordinator/query_coordinator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -322,11 +322,13 @@ impl<'a> StageCoordinator<'a> {
// metrics collection process might outlive the query's lifetime.
#[allow(clippy::disallowed_methods)]
tokio::spawn(async move {
let mut received_metrics = false;
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);
received_metrics = true;
}
Comment thread
gabotechs marked this conversation as resolved.
Outdated
}
WorkerToCoordinatorMsg::LoadInfo(load_info) => {
Expand All @@ -347,6 +349,18 @@ impl<'a> StageCoordinator<'a> {
}
}
}
if !received_metrics {
if let Some(task_metrics) = task_metrics {
// An unexecuted task sends no metrics; still complete its wait.
task_metrics.insert(
task_key,
TaskMetrics {
pre_order_plan_metrics: vec![],
task_metrics: Default::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
1 change: 0 additions & 1 deletion tests/metrics_collection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -382,7 +382,6 @@ 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;
Expand Down
Loading