Metrics#
DataFusion operators expose runtime metrics so you can understand where time is spent and how much data flows through the pipeline. See more in EXPLAIN ANALYZE.
Common Metrics#
BaselineMetrics#
BaselineMetrics are available in most physical operators to capture common measurements.
Metric |
Description |
|---|---|
elapsed_compute |
CPU time the operator actively spends processing work. |
output_rows |
Total number of rows the operator produces. |
output_bytes |
Memory usage of all output batches. Note: This value may be overestimated. If multiple output |
output_batches |
Total number of output batches the operator produces. |
Operator-specific Metrics#
FilterExec#
Metric |
Description |
|---|---|
selectivity |
Selectivity of the filter, calculated as output_rows / input_rows |
HashJoinExec#
HashJoinExec also exposes the common BaselineMetrics. Its
elapsed_compute metric is the sum of the build-side collection time and the
subsequent join processing time.
Metric |
Description |
|---|---|
build_time |
Total time spent collecting and building the build side of the join. |
build_input_batches |
Number of input batches consumed from the build side. |
build_input_rows |
Number of input rows consumed from the build side. |
build_mem_used |
Peak tracked memory used by the build side, in bytes. |
join_time |
Total time spent processing the join after build-side collection. |
input_batches |
Number of input batches consumed from the probe side. |
input_rows |
Number of input rows consumed from the probe side. |
probe_hit_rate |
Fraction of probe-side rows with a build-side join-key match before applying any join filter. |
avg_fanout |
Average number of build-side join-key matches per matched probe-side row before applying any join filter. |
array_map_created_count |
Number of times |
AggregateExec#
AggregateExec exposes the common BaselineMetrics, the operator-level
metrics below, and a timer per aggregate expression and execution phase.
Metric |
Description |
|---|---|
time_calculating_group_ids |
Time spent preparing group keys; see the note below for path-specific coverage. |
aggregate_arguments_time |
Total time spent evaluating the inputs to the aggregate functions. |
aggregation_time |
Time spent invoking accumulator |
emitting_time |
Time spent materializing group values and accumulator |
topk_maintenance_time |
Time spent maintaining Grouped TopK’s priority map, including batch setup, insertion, comparison, and NULL handling. |
skipped_aggregation_rows |
Number of input rows passed through without aggregating them, when partial aggregation is skipped. |
reduction_factor |
Rows emitted per row consumed by a partial aggregation, displayed as |
spill_count |
Number of spill files written when the aggregation exceeds its memory budget. |
spilled_bytes |
Total number of bytes written to spill files. |
spilled_rows |
Total number of rows written to spill files. |
peak_mem_used |
Peak tracked memory held by the grouped aggregation, in bytes; recorded by the fallback grouped hash path only. |
These operator-level metrics are recorded by the grouped aggregation paths
only: an AggregateExec without a GROUP BY reports just BaselineMetrics
and the per-aggregate timers. reduction_factor and skipped_aggregation_rows
are recorded in partial mode only. skipped_aggregation_rows is recorded only
when partial-aggregation skipping is enabled for a single, non-grouping-sets
GROUP BY whose input is not ordered by its grouping expressions; a
datafusion.execution.skip_partial_aggregation_probe_ratio_threshold of >= 1.0
disables the feature.
time_calculating_group_ids covers both grouping-expression evaluation and
resolving the resulting rows to group IDs, including interning and ordering
setup. aggregation_time covers only accumulator update and merge calls.
On accumulator-backed grouped aggregation paths, state and evaluate are
included in emitting_time, including when partial aggregation materializes
state under memory pressure. For an AggregateExec without a GROUP BY, the
per-aggregate state and evaluate timers are recorded during
elapsed_compute; it does not report grouped operator metrics such as
emitting_time. Partial aggregation that skips aggregation reports
convert_to_state instead of aggregation_time for those rows.
When the specialized Grouped TopK path is selected for a limited grouped
aggregate, it has no accumulators. Its group-key expression evaluation is
included in time_calculating_group_ids, and its priority-map work is reported
by topk_maintenance_time; it does not report accumulator phases. Planner
selection is query-shape dependent: a limited DISTINCT query without
ordering requirements uses the regular aggregate path instead.
The per-aggregate timers are named agg_expr_{index}_{phase}_time, where
index is the zero-based position of an aggregate expression in the operator
and phase is one of the following:
Phase |
Description |
|---|---|
|
Evaluating the aggregate’s argument expressions into input arrays. |
|
Updating an accumulator from raw input values. |
|
Merging partial accumulator states. |
|
Obtaining an accumulator’s intermediate state for partial output or aggregate spilling. |
|
Converting raw aggregate inputs directly to partial state without normal accumulator updates. |
|
Evaluating an accumulator to its final result. |
For example, when a partial stage is planned, the partial AggregateExec for
SELECT SUM(a), SUM(b) FROM t reports agg_expr_0_arguments_time and
agg_expr_0_update_time for SUM(a),
and agg_expr_1_arguments_time and agg_expr_1_update_time for SUM(b). The
index is positional and refers to the same position in the aggr=[...] list
printed on the operator’s plan line, which is how an indexed timer is mapped
back to an aggregate expression. Because the index is part of the metric name,
otherwise identical functions over different columns stay distinct when
per-partition metrics are combined.
Each per-aggregate timer additionally carries an aggregate label holding the
rendered aggregate expression (for example, sum(t.a)). Combining metrics
across partitions drops labels, so this label is only shown in the “Plan with
Full Metrics” section of EXPLAIN ANALYZE VERBOSE, which reports metrics per
partition.
arguments is recorded in every mode. The accumulator phases that are present
depend on the aggregate mode and implementation. For non-grouped aggregation,
partial mode records update and state, partial reduce mode records merge
and state, final mode records merge and evaluate, and single mode records
update and evaluate. Hash aggregation uses update, state, and
convert_to_state in partial mode; merge and state in partial-reduce mode;
merge, state, and evaluate in final mode; and update, state, merge,
and evaluate in single mode. Its state timers measure intermediate-state
emission, including during spilling. The grouped TopK aggregate path records
only the per-aggregate arguments timer, because it maintains values directly
rather than using accumulators.
Where an aggregate has a FILTER clause, its evaluation is included in
aggregate_arguments_time. Per-aggregate arguments timers include the filter
when it is evaluated per aggregate. The legacy grouped hash path evaluates
filters collectively, so its per-aggregate arguments timers cover argument
expressions only; their sum need not equal aggregate_arguments_time.
Aggregate implementations can also expose optional internal submetrics. These
use the agg_expr_{index}_internal_{subphase}_time naming and aggregate
label. The internal segment keeps them separate from call-boundary timers;
subphase is a stable identifier owned and documented by the aggregate
implementation. They are registered lazily only when an aggregate requests
them, so aggregates without internal submetrics add no metrics. Registration
is per (aggregate expression index, subphase, partition): replacement
accumulators in that partition share the same time, and normal metric display
combines that time across partitions. An aggregate may request its submetric
during accumulator construction, so it can appear even when its input is empty.
For example, array_agg(DISTINCT ...) records the time spent deduplicating
input values as agg_expr_{index}_internal_distinct_time. Grouped accumulation
records this once per input batch, rather than once per group, to avoid making
metric collection proportional to group cardinality. These submetrics complement
the update, merge, state, and evaluate timers rather than subdividing or
replacing them. An internal submetric may therefore overlap its enclosing phase
timer; it is a supplementary diagnostic and must not be added to phase timings
as a breakdown.
Except for the Summary metric reduction_factor, these operator-level and
per-aggregate metrics are Dev metrics. They appear in EXPLAIN ANALYZE when
datafusion.explain.analyze_level includes Dev (the default), but are omitted
at the Summary level. The normal display combines partitions; use EXPLAIN ANALYZE VERBOSE to additionally show the per-partition values together with
each per-aggregate timer’s aggregate label. For a query
such as the following, the per-expression metrics stay readable, and the
operator’s aggr=[...] list names the aggregate behind each timer index:
EXPLAIN ANALYZE
SELECT k, SUM(a), SUM(b), COUNT(c)
FROM t
GROUP BY k;
The abbreviated partial-aggregate plan line below shows that mapping (timings and unrelated metrics are omitted):
AggregateExec: mode=Partial, gby=[k@0 as k], aggr=[sum(t.a), sum(t.b), count(t.c)],
metrics=[..., agg_expr_0_arguments_time=..., agg_expr_1_arguments_time=...,
agg_expr_2_arguments_time=...]
Thus agg_expr_0_* is sum(t.a), agg_expr_1_* is sum(t.b), and
agg_expr_2_* is count(t.c). In verbose output, the corresponding
aggregate labels provide the same mapping per partition.
TODO#
Add metrics for the remaining operators