alexy opened a new issue, #26071: URL: https://github.com/apache/datafusion/issues/26071
### Is your feature request related to a problem or challenge? Every operator's `BaselineMetrics` records `output_bytes` for every output batch through `RecordOutput::record_output` (`datafusion/physical-expr-common/src/metrics/baseline.rs`), which calls `get_record_batch_memory_size`. That function converts each column, and each child of a nested column, to `ArrayData` and inserts every buffer's address into a `HashSet`, so that buffers shared between columns are counted once. Exact accounting matters for memory reservations; for a per-batch metric it is a cost every query pays, and it grows with the width and nesting of the schema rather than with the data. Measured on `main` (c0e872f), release build, Apple M1 Max, one batch of `Int32` struct columns: | batch | `get_record_batch_memory_size` | sum of `Array::get_array_memory_size` | |---|---:|---:| | 35 struct columns x 20 fields, 100 rows | 21.9 µs | 1.6 µs | | 35 struct columns x 20 fields, 4,000 rows | 15.1 µs | 1.5 µs | | 5 struct columns x 5 fields, 4,000 rows | 0.58 µs | 0.05 µs | That is per operator per batch. We hit it running one wide plan repeatedly (a game's world, 35 struct columns, a few thousand rows, a few hundred operators, through [Sail](https://github.com/lakehq/sail) on DataFusion 55.1): in a sampled profile, `BaselineMetrics::record_poll` -> `get_record_batch_memory_size` -> `StructArray::to_data` was about 12% of executing the plan, and computing the metric as below took a run from 40 ms to 34 ms. ### Describe the solution you'd like Compute `output_bytes` without `ArrayData`, for example: ```rust fn output_bytes(batch: &RecordBatch) -> usize { batch.columns().iter().map(|c| c.get_array_memory_size()).sum() } ``` This can count a buffer shared between columns more than once (in the table above it reports 350,840 bytes against 280,000 for the 100-row batch, which shares nothing; the difference is the arrays' own structures), which seems acceptable for a metric; `get_record_batch_memory_size` would stay for memory accounting. Alternatively, `output_bytes` could be computed only when a metrics level that shows it is requested. ### Describe alternatives you've considered Making `get_record_batch_memory_size` itself cheaper for nested arrays (walking the arrays' buffers without building `ArrayData`) would help memory accounting too, but keeps the hash set per batch. ### Additional context Related: #26065 (struct literal comparisons through `ArrayData` in physical planning) and the planning-speed work in #19795. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
