Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 25 additions & 8 deletions datafusion/physical-expr-common/src/metrics/baseline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@
use std::{borrow::Cow, collections::BTreeMap, sync::Arc, task::Poll};

use arrow::record_batch::RecordBatch;
use datafusion_common::{Result, utils::memory::get_record_batch_memory_size};
use datafusion_common::Result;

use super::{
Count, ExecutionPlanMetricsSet, Metric, MetricBuilder, MetricsSet, Time, Timestamp,
Expand Down Expand Up @@ -62,9 +62,14 @@ pub struct BaselineMetrics {

/// Memory usage of all output batches.
///
/// Note: This value may be overestimated. If multiple output `RecordBatch`
/// instances share underlying memory buffers, their sizes will be counted
/// multiple times.
/// Computed as the sum of each output column's `Array::get_array_memory_size`
/// so that recording this metric never builds `ArrayData`. Use
/// `datafusion_common::utils::memory::get_record_batch_memory_size` instead
/// when exact, buffer-deduplicated accounting is required.
///
/// Note: This value may be overestimated. A buffer shared between two
/// columns, or between an array and its children, is counted once per
/// reference, and each array's own in-memory structure is included.
/// Issue: <https://github.com/apache/datafusion/issues/16841>
output_bytes: Count,

Expand Down Expand Up @@ -312,6 +317,20 @@ impl SplitMetrics {
}
}

/// Returns the total memory of `batch`'s arrays, computed without building
/// `ArrayData`.
///
/// Unlike `datafusion_common::utils::memory::get_record_batch_memory_size`,
/// this does not deduplicate buffers shared between columns, and includes each
/// `Array`'s own in-memory structure.
fn output_bytes(batch: &RecordBatch) -> usize {
batch
.columns()
.iter()
.map(|column| column.get_array_memory_size())
.sum()
}

/// Trait for things that produce output rows as a result of execution.
pub trait RecordOutput {
/// Record that some number of output rows have been produced
Expand All @@ -331,8 +350,7 @@ impl RecordOutput for usize {
impl RecordOutput for RecordBatch {
fn record_output(self, bm: &BaselineMetrics) -> Self {
bm.record_output(self.num_rows());
let n_bytes = get_record_batch_memory_size(&self);
bm.output_bytes.add(n_bytes);
bm.output_bytes.add(output_bytes(&self));
bm.output_batches.add(1);
self
}
Expand All @@ -341,8 +359,7 @@ impl RecordOutput for RecordBatch {
impl RecordOutput for &RecordBatch {
fn record_output(self, bm: &BaselineMetrics) -> Self {
bm.record_output(self.num_rows());
let n_bytes = get_record_batch_memory_size(self);
bm.output_bytes.add(n_bytes);
bm.output_bytes.add(output_bytes(self));
bm.output_batches.add(1);
self
}
Expand Down
12 changes: 6 additions & 6 deletions docs/source/user-guide/metrics.md
Original file line number Diff line number Diff line change
Expand Up @@ -27,12 +27,12 @@ DataFusion operators expose runtime metrics so you can understand where time is

`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 `RecordBatch` instances share underlying memory buffers, their sizes will be counted multiple times. |
| output_batches | Total number of output batches the operator produces. |
| 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, computed as the sum of each column's `Array::get_array_memory_size`. Note: This value may be overestimated: a buffer shared between two columns, or between an array and its children, is counted once per reference, and each array's own in-memory structure is included. |
| output_batches | Total number of output batches the operator produces. |

## Operator-specific Metrics

Expand Down