Repository navigation
perf: avoid ArrayData conversion in BaselineMetrics for every column - #26105
Open
minhh-nguyen wants to merge 2 commits into
Open
minhh-nguyen wants to merge 2 commits into
minhh-nguyen wants to merge 2 commits into
Conversation
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
force-pushed
the
perf/26071-avoid-arraydata-conversion
branch
from
October 7, 2026 08:58
209816f to
e2014b3
Compare
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Which issue does this PR close?
Rationale for this change
EXPLAIN ANALYZEreportsoutput_bytesfor every operator, and recording that metric was costing more than the metric is worth.BaselineMetricscalledget_record_batch_memory_sizeon each output batch, which walksevery 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_sizewas about 12% of planexecution.
This PR computes the metric without materializing
ArrayData.What changes are included in this PR?
BaselineMetricsrecordsoutput_bytesas the sum of each output column'sArray::get_array_memory_sizeinstead of callingget_record_batch_memory_size.get_record_batch_memory_sizeandRecordBatchMemoryCounterare unchanged, so memory reservations keep exact, deduplicated accounting.Every operator that reserves against batch memory still uses them directly.
output_bytesis now an approximation, and update the metric description indocs/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 ANALYZEnow reportsoutput_bytesas 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 eacharray's own in-memory struct is included. The metric description in
docs/source/user-guide/metrics.mdwas updated to say so. No query results, public APIs, or memory accounting behavior change.