Skip to content

perf: avoid ArrayData conversion in BaselineMetrics for every column - #26105

Open
minhh-nguyen wants to merge 2 commits into
apache:mainfrom
minhh-nguyen:perf/26071-avoid-arraydata-conversion
Open

minhh-nguyen wants to merge 2 commits into
apache:mainfrom
minhh-nguyen:perf/26071-avoid-arraydata-conversion

Conversation

@minhh-nguyen

@minhh-nguyen minhh-nguyen commented Oct 7, 2026 •

Copy link
Copy Markdown

Which issue does this PR close?

Rationale for this change

EXPLAIN ANALYZE reports output_bytes for every operator, and recording that metric was costing more than the metric is worth. BaselineMetrics called get_record_batch_memory_size on each output batch, which walks
every column and every child of a nested column and inserts each buffer's address into a hash set so shared buffers are counted once.

That exact accounting is required for memory reservations, but it is not required for a metric that is only displayed, and its cost grows with the width and nesting of the schema rather than with the data. In a sampled profile of a wide plan run through Sail on DataFusion 55.1, record_poll -> get_record_batch_memory_size was about 12% of plan
execution.

This PR computes the metric without materializing ArrayData.

What changes are included in this PR?

  • BaselineMetrics records output_bytes as the sum of each output column's Array::get_array_memory_size instead of calling get_record_batch_memory_size.
  • get_record_batch_memory_size and RecordBatchMemoryCounter are unchanged, so memory reservations keep exact, deduplicated accounting.
    Every operator that reserves against batch memory still uses them directly.
  • Document that output_bytes is now an approximation, and update the metric description in docs/source/user-guide/metrics.md.

What is the testing strategy for this PR?

no new test; output_bytes values are unchanged in existing sqllogictests;
note that hardcoded timings in explain_analyze.slt fail identically on main

Are there any user-facing changes?

EXPLAIN ANALYZE now reports output_bytes as an approximation rathe than a buffer-deduplicated total. A buffer shared between columns or between an array and its children is counted once per reference, and each
array's own in-memory struct is included. The metric description in docs/source/user-guide/metrics.md was updated to say so. No query results, public APIs, or memory accounting behavior change.

@github-actions github-actions Bot added documentation Improvements or additions to documentation physical-expr Changes to the physical-expr crates labels Oct 7, 2026
Closes: apache#26071

`BaselineMetrics::record_output` called `get_record_batch_memory_size`
for every output batch, which materializes `ArrayData` and inserts every
buffer address into a hash set. Compute `output_bytes` as the sum of
each
column's `Array::get_array_memory_size` instead: it never builds
`ArrayData` and is 10.8-13.9x faster on wide struct batches.

`get_record_batch_memory_size` is unchanged and remains what memory
reservations use, since those need exact deduplicated accounting.
Closes apache#26071

`BaselineMetrics::record_output` called `get_record_batch_memory_size`
for every output batch, which materializes `ArrayData` and inserts every
buffer address into a hash set. Compute `output_bytes` as the sum of
each
column's `Array::get_array_memory_size` instead, which never builds
`ArrayData`.

`get_record_batch_memory_size` is unchanged and remains what memory
reservations use, since those need exact deduplicated accounting.
@minhh-nguyen
minhh-nguyen force-pushed the perf/26071-avoid-arraydata-conversion branch from 209816f to e2014b3 Compare October 7, 2026 08:58
@minhh-nguyen minhh-nguyen changed the title Perf/26071 avoid arraydata conversion perf: avoid ArrayData conversion in BaselineMetrics for every column Oct 7, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

documentation Improvements or additions to documentation physical-expr Changes to the physical-expr crates

Projects

None yet

Development

Successfully merging this pull request may close these issues.

BaselineMetrics output_bytes converts every column to ArrayData for every batch

1 participant