Skip to content

fix: finalize metrics for unexecuted distributed tasks - #744

Closed
sesteves wants to merge 2 commits into
datafusion-contrib:mainfrom
sesteves:fix/metrics-terminal-task-state
Closed

sesteves wants to merge 2 commits into
datafusion-contrib:mainfrom
sesteves:fix/metrics-terminal-task-state

Conversation

@sesteves

@sesteves sesteves commented Sep 23, 2026 •

Copy link
Copy Markdown
Contributor

Summary

  • Finalize metrics for planned tasks that never receive ExecuteTask, preserving actual AQE sampling I/O and reporting skipped static tasks without fabricated execution metrics.
  • Distinguish pending tasks from channels closed without a report, and bound coordinator cleanup, metrics, and completed dynamic-filter waits with one configurable post-query deadline.
  • Document the finalization deadline and cover AQE and static short-circuiting with regression tests.

Closes #739

Benchmark

  • Benchmarked the implementation against the baseline using paired release runs of the 22-query TPC-H SF1 suite. Timings changed direction across runs, with no consistent difference. The later test/documentation simplifications and small worker reporting refactor were not rebenchmarked.

@gabotechs

Copy link
Copy Markdown
Collaborator

Hi @sesteves! thanks for the PR!

I think it would make sense to split this in two, it'd be awesome if we could first ship a fix for #739 in isolation. I think the fix for that can be quite scoped.

#685 will require some alignment and some design, and I'd not want that to block a fix for #739.

What do you think?

@sesteves

Copy link
Copy Markdown
Contributor Author

Agreed — I’ve updated this PR to focus on #739. I’ll handle #685 separately after we align on the API.

@gabotechs

Copy link
Copy Markdown
Collaborator

Nice! I was expecting the solution to be way more self contained, I was not expecting more than 10 or 20 lines of code for this fix, seems like we have close to 800.

Your integration tests as good as they are, are you fine if we make a first PR adding the failing tests (marked as #[ignore]-ed for now), and we iterate on a fix on a follow up?

.await??;
let rewritten = rewrite_with_metrics_within_timeout(with_filters).await?;
let display = display_plan_ascii(rewritten.as_ref(), true);
assert_contains!(&display, "task_not_executed=");

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🤔 I'm not sure if we'd want to surface these. In the same way that datafusion does not expose a paritions_not_executed metric, it might also make sense to follow the same approach for remote tasks.

I'd expect a not executed task to just be missing its entry in the metrics of a node, something like this:

    ┌───── Stage 2 ── tasks=3, partitions=3 cpu_cost=18.6 KB, estimated_output_bytes=7.4 KB, estimated_pct_sampled=100, memory_cost=0.0 B, network_cost=0.0 B, plan_added_at={1:3.73ms}, plan_executed_at={1:22.78ms}, plan_finished_at={1:25.85ms}
    │ RepartitionExec: partitioning=Hash([RainToday@1], 3), input_partitions=3, metrics=[spill_count={1:0}, spilled_bytes={1:0.0 B}, spilled_rows={1:0}, fetch_time={1:3.92µs}, repartition_time={1:3ns}, send_time={1:9ns}]
    │   SamplerExec: partitions=3, metrics=[elapsed_compute={1:1ns}, max_mem_used={1:0.0 B},

Where only task with index 1 out of 0,1,2 where executed.

@sesteves sesteves closed this Sep 25, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Metrics finalization can hang on planned but never executed tasks with AQE or static planning

2 participants