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
5 changes: 0 additions & 5 deletions datafusion/physical-plan/src/aggregates/order/full.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,6 @@
// under the License.

use datafusion_expr::EmitTo;
use std::mem::size_of;

/// Tracks grouping state when the data is ordered entirely by its
/// group keys
Expand Down Expand Up @@ -143,10 +142,6 @@ impl GroupOrderingFull {
}
};
}

pub(crate) fn size(&self) -> usize {
size_of::<Self>()
}
}

impl Default for GroupOrderingFull {
Expand Down
46 changes: 42 additions & 4 deletions datafusion/physical-plan/src/aggregates/order/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -148,12 +148,15 @@ impl GroupOrdering {
}

/// Returns the size of memory used by the ordering state, in bytes.
///
/// Includes the enum descriptor once (covering the inline active variant)
/// and any heap allocations retained by that variant.
pub fn size(&self) -> usize {
size_of::<Self>()
+ match self {
GroupOrdering::None => 0,
GroupOrdering::Partial(partial) => partial.size(),
GroupOrdering::Full(full) => full.size(),
// These variants hold only inline state, already counted above.
GroupOrdering::None | GroupOrdering::Full(_) => 0,
GroupOrdering::Partial(partial) => partial.heap_size(),
}
}
}
Expand All @@ -164,7 +167,42 @@ mod tests {

use std::sync::Arc;

use arrow::array::Int32Array;
use arrow::array::{Int32Array, StringArray};
use datafusion_common::ScalarValue;

#[test]
fn test_size_inline_only() -> Result<()> {
let expected = size_of::<GroupOrdering>();
assert_eq!(GroupOrdering::None.size(), expected);
let ordering = GroupOrdering::try_new(&InputOrderMode::Sorted)?;
assert_eq!(ordering.size(), expected);
Ok(())
}
Comment on lines +173 to +180

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

One test is enough here: assert that None and a Sorted ordering both report exactly size_of::<GroupOrdering>(). None | Full(_) never reads the variant, so the new_groups, remove_groups, input_done and reset steps cannot change the result. The first Full assertion already catches the old code, and test_size_none passed before this PR too.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I’ll combine the None and Sorted assertions into one descriptor-size test; the Sorted assertion catches the original double charge.


#[test]
fn test_size_partial_retained_allocations() -> Result<()> {
// The direct partial constructor retains spare order-index capacity.
let mut order_indices = Vec::with_capacity(32);
order_indices.push(0);
let expected =
size_of::<GroupOrdering>() + order_indices.capacity() * size_of::<usize>();
let mut ordering =
GroupOrdering::Partial(GroupOrderingPartial::try_new(order_indices)?);
assert_eq!(ordering.size(), expected);

// A variable-width key checks both the scalar descriptor and its payload.
let key = "retained sort key";
let batch_group_values: Vec<ArrayRef> =
vec![Arc::new(StringArray::from(vec!["a", key]))];
ordering.new_groups(&batch_group_values, &[0, 1], 2)?;
let in_progress = expected + ScalarValue::Utf8(Some(key.to_owned())).size();
assert_eq!(ordering.size(), in_progress);

// Completing drops the key, but retains order-index capacity.
ordering.input_done();
assert_eq!(ordering.size(), expected);
Ok(())
}

#[test]
fn test_oom_emit_to_none_ordering() {
Expand Down
10 changes: 5 additions & 5 deletions datafusion/physical-plan/src/aggregates/order/partial.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,6 @@
// under the License.

use std::cmp::Ordering;
use std::mem::size_of;
use std::sync::Arc;

use arrow::array::ArrayRef;
Expand Down Expand Up @@ -102,7 +101,7 @@ enum State {
}

impl State {
fn size(&self) -> usize {
fn heap_size(&self) -> usize {
match self {
State::Taken => 0,
State::Start => 0,
Expand Down Expand Up @@ -267,9 +266,10 @@ impl GroupOrderingPartial {
Ok(())
}

/// Return the size of memory allocated by this structure
pub(crate) fn size(&self) -> usize {
size_of::<Self>() + self.order_indices.allocated_size() + self.state.size()
/// Returns retained heap allocations, excluding the inline descriptor
/// already counted by [`super::GroupOrdering::size`].
pub(crate) fn heap_size(&self) -> usize {
self.order_indices.allocated_size() + self.state.heap_size()
}
}

Expand Down
Loading