Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
18 commits
Select commit Hold shift + click to select a range
70dee35
fix: deduplicate StringView/BinaryView buffer refs in CollectLeft Has…
mohitgurav20 Sep 24, 2026
19dfb9c
Merge branch 'main' into fix/dedup-stringview-buffers-collectleft-joins
mohitgurav20 Sep 25, 2026
c9a9fac
Merge branch 'main' into fix/dedup-stringview-buffers-collectleft-joins
mohitgurav20 Sep 26, 2026
2df2dc3
Merge branch 'main' into fix/dedup-stringview-buffers-collectleft-joins
mohitgurav20 Sep 27, 2026
de9d639
style: format code
mohitgurav20 Sep 25, 2026
32e95d5
fix: convert unique_buffers to Arc<[Buffer]>
mohitgurav20 Sep 27, 2026
7ed87d1
Merge branch 'main' into fix/dedup-stringview-buffers-collectleft-joins
mohitgurav20 Sep 27, 2026
905f0b3
style: run cargo fmt and fix clippy warnings
mohitgurav20 Sep 27, 2026
3995a49
style: run cargo fmt
mohitgurav20 Sep 28, 2026
8865702
chore: restore .gitignore to main (remove accidental UTF-16 encoding)
mohitgurav20 Oct 6, 2026
7a359a8
Merge branch 'main' into fix/dedup-stringview-buffers-collectleft-joins
mohitgurav20 Oct 6, 2026
3500877
test: add unit tests for deduplicate_view_array_buffers and deduplica…
mohitgurav20 Oct 6, 2026
980a2b8
Merge branch 'fix/dedup-stringview-buffers-collectleft-joins' of http…
mohitgurav20 Oct 6, 2026
16ed8a1
test: improve coverage for view buffer deduplication paths
mohitgurav20 Oct 6, 2026
42e3b02
test: eliminate Result<?> error branches and cover retained>held grow…
mohitgurav20 Oct 6, 2026
41d84a9
fix: make deduplicate_record_batch_view_buffers infallible to fix cov…
mohitgurav20 Oct 6, 2026
e5eb28d
fix: use HashTable::allocation_size() in GroupValuesColumn::size()
mohitgurav20 Oct 6, 2026
3e20bd5
Merge branch 'main' into fix/exact-hashtable-accounting-group-values-…
mohitgurav20 Oct 6, 2026
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
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ pub(super) use ordered::GroupValuesOrdered;
pub mod primitive;
pub mod row_backed;

use std::mem::{self, size_of};
use std::mem;
use std::sync::Arc;

use crate::aggregates::group_values::GroupValues;
Expand All @@ -49,7 +49,7 @@ use datafusion_common::hash_utils::RandomState;
use datafusion_common::hash_utils::create_hashes;
use datafusion_common::utils::{has_float_leaf, normalize_float_zero};
use datafusion_common::{Result, not_impl_err};
use datafusion_execution::memory_pool::proxy::{HashTableAllocExt, VecAllocExt};
use datafusion_execution::memory_pool::proxy::VecAllocExt;
use datafusion_expr::{EmitTo, GroupSelection};
use datafusion_physical_expr::binary_map::OutputType;

Expand Down Expand Up @@ -194,9 +194,6 @@ pub struct GroupValuesColumn<const STREAMING: bool> {
///
map: HashTable<(u64, GroupIndexView)>,

/// The size of `map` in bytes
map_size: usize,

/// The lists for group indices with the same hash value
///
/// It is possible that hash value collision exists,
Expand Down Expand Up @@ -290,7 +287,6 @@ impl<const STREAMING: bool> GroupValuesColumn<STREAMING> {
group_index_lists: Vec::new(),
emit_group_index_list_buffer: Vec::new(),
vectorized_operation_buffers: VectorizedOperationBuffers::default(),
map_size: 0,
group_values,
hashes_buffer: Default::default(),
random_state: crate::aggregates::AGGREGATION_HASH_SEED,
Expand Down Expand Up @@ -431,10 +427,10 @@ impl<const STREAMING: bool> GroupValuesColumn<STREAMING> {
}

// for hasher function, use precomputed hash value
self.map.insert_accounted(
self.map.insert_unique(
target_hash,
(target_hash, GroupIndexView::new_inlined(group_idx as u64)),
|(hash, _group_index)| *hash,
&mut self.map_size,
);
group_idx
}
Expand Down Expand Up @@ -551,10 +547,10 @@ impl<const STREAMING: bool> GroupValuesColumn<STREAMING> {

// Insert the `group index view` and its hash into `map`
// for hasher function, use precomputed hash value
self.map.insert_accounted(
self.map.insert_unique(
target_hash,
(target_hash, group_index_view),
|(hash, _)| *hash,
&mut self.map_size,
);

// Add row index to `vectorized_append_row_indices`
Expand Down Expand Up @@ -1068,7 +1064,12 @@ impl<const STREAMING: bool> GroupValues for GroupValuesColumn<STREAMING> {

fn size(&self) -> usize {
let group_values_size: usize = self.group_values.iter().map(|v| v.size()).sum();
group_values_size + self.map_size + self.hashes_buffer.allocated_size()
// `HashTable::allocation_size` reports the complete retained allocation —
// including hashbrown control bytes and trailing layout — which is exactly
// what we want here. This follows the same approach as `ArrowBytesMap::size()`.
group_values_size
+ self.map.allocation_size()
+ self.hashes_buffer.allocated_size()
}

fn is_empty(&self) -> bool {
Expand Down Expand Up @@ -1210,7 +1211,6 @@ impl<const STREAMING: bool> GroupValues for GroupValuesColumn<STREAMING> {
.expect("schema previously validated in try_new");
self.map.clear();
self.map.shrink_to(num_rows, |_| 0); // hasher does not matter since the map is cleared
self.map_size = self.map.capacity() * size_of::<(u64, usize)>();
self.hashes_buffer.clear();
self.hashes_buffer.shrink_to(num_rows);

Expand Down Expand Up @@ -1258,8 +1258,8 @@ mod tests {
compute::{concat_batches, take},
util::pretty::pretty_format_batches,
};
use datafusion_common::utils::proxy::HashTableAllocExt;
use datafusion_expr::{EmitTo, GroupSelection};
use std::mem::size_of;

use crate::aggregates::group_values::{
GroupValues, multi_group_by::GroupValuesColumn,
Expand Down Expand Up @@ -2655,10 +2655,10 @@ mod tests {
group_index: u64,
) {
let group_index_view = GroupIndexView::new_inlined(group_index);
group_values.map.insert_accounted(
group_values.map.insert_unique(
hash_key,
(hash_key, group_index_view),
|(hash, _)| *hash,
&mut group_values.map_size,
);
}

Expand All @@ -2670,10 +2670,159 @@ mod tests {
let list_offset = group_values.group_index_lists.len();
let group_index_view = GroupIndexView::new_non_inlined(list_offset as u64);
group_values.group_index_lists.push(group_indices);
group_values.map.insert_accounted(
group_values.map.insert_unique(
hash_key,
(hash_key, group_index_view),
|(hash, _)| *hash,
&mut group_values.map_size,
);
}

// -----------------------------------------------------------------------
// Tests for exact hash-table allocation accounting (issue #25736)
// -----------------------------------------------------------------------

/// After enough groups are inserted to force a hashbrown table growth,
/// `size()` must include the full `allocation_size()` of the map — control
/// bytes and trailing layout included — not just an entry-capacity estimate.
#[test]
fn map_allocation_size_included_in_size_after_growth() {
let schema: SchemaRef =
Arc::new(Schema::new(vec![Field::new("k", DataType::Int64, false)]));
let mut gv: GroupValuesColumn<false> =
GroupValuesColumn::try_new(Arc::clone(&schema)).unwrap();

// Insert enough distinct rows to trigger at least one table resize.
let n: usize = 128;
let keys: Vec<i64> = (0..n as i64).collect();
let col = Arc::new(Int64Array::from(keys)) as ArrayRef;
let mut groups = vec![];
gv.intern(&[col], &mut groups).unwrap();

let reported = gv.size();
let map_alloc = gv.map.allocation_size();

// The reported size must be at least as large as the raw map allocation.
assert!(
reported >= map_alloc,
"size() ({reported}) must be >= map.allocation_size() ({map_alloc})"
);

// And the map must actually hold a non-trivial allocation after growth.
assert!(
map_alloc > 0,
"map.allocation_size() should be > 0 after inserting {n} groups"
);
}

/// When hash collisions are forced, the collision-chain list allocation
/// (`group_index_lists`) must be accounted for independently of the
/// hash-table term. The two must not be conflated.
#[test]
fn collision_chain_size_is_separate_from_map_allocation() {
let schema: SchemaRef =
Arc::new(Schema::new(vec![Field::new("k", DataType::Int64, false)]));
let mut gv: GroupValuesColumn<false> =
GroupValuesColumn::try_new(Arc::clone(&schema)).unwrap();

// Manually insert two entries with the same hash key to force a
// collision chain (non-inlined GroupIndexView).
insert_inline_group_index_view(&mut gv, 42, 0);
insert_non_inline_group_index_view(&mut gv, 42, vec![0, 1]);

let map_alloc = gv.map.allocation_size();
let chain_size: usize = gv
.group_index_lists
.iter()
.map(|list| list.capacity() * size_of::<usize>())
.sum();

// Both terms must be individually non-zero.
assert!(map_alloc > 0, "map allocation must be > 0");
assert!(chain_size > 0, "collision-chain allocation must be > 0");

// The reported total must include both independently.
let reported = gv.size();
assert!(
reported >= map_alloc,
"size() must cover the map allocation"
);
}

/// After a partial emit (`EmitTo::First`), the map retains its allocated
/// capacity even though logical entries were removed. `size()` must still
/// reflect the retained allocation, not drop to zero.
#[test]
fn size_reflects_retained_capacity_after_partial_emit() {
let schema: SchemaRef =
Arc::new(Schema::new(vec![Field::new("k", DataType::Int64, false)]));
let mut gv: GroupValuesColumn<false> =
GroupValuesColumn::try_new(Arc::clone(&schema)).unwrap();

let keys: Vec<i64> = (0..64i64).collect();
let col = Arc::new(Int64Array::from(keys)) as ArrayRef;
let mut groups = vec![];
gv.intern(&[col], &mut groups).unwrap();

let size_before = gv.size();
let map_alloc_before = gv.map.allocation_size();

// Emit only the first 32 groups; the table allocation should be retained.
gv.emit(EmitTo::First(32)).unwrap();

let size_after = gv.size();
let map_alloc_after = gv.map.allocation_size();

// Capacity is not released by a partial emit — allocation should persist.
assert!(
map_alloc_after > 0,
"map allocation should be retained after partial emit, got {map_alloc_after}"
);
assert!(
size_after > 0,
"size() should remain > 0 after partial emit, got {size_after}"
);

// The allocation reported before and after should both be real and
// consistent with what the map actually holds.
assert_eq!(
map_alloc_before, map_alloc_after,
"partial emit must not shrink the map allocation (before={map_alloc_before}, after={map_alloc_after})"
);
let _ = size_before; // used implicitly via the assertions above
}

/// After a full emit followed by re-use, `size()` must still produce
/// correct group values and the map allocation is reported accurately.
#[test]
fn reuse_after_full_emit_produces_correct_groups() {
let schema: SchemaRef =
Arc::new(Schema::new(vec![Field::new("k", DataType::Int64, false)]));
let mut gv: GroupValuesColumn<false> =
GroupValuesColumn::try_new(Arc::clone(&schema)).unwrap();

// First pass: intern 4 distinct keys.
let first_keys: Vec<i64> = vec![10, 20, 30, 40];
let col1 = Arc::new(Int64Array::from(first_keys.clone())) as ArrayRef;
let mut groups = vec![];
gv.intern(&[col1], &mut groups).unwrap();
assert_eq!(groups, vec![0, 1, 2, 3]);

// Full emit clears the map.
let emitted = gv.emit(EmitTo::All).unwrap();
assert_eq!(emitted.len(), 1);

// Second pass after reuse: intern 3 different keys, numbering restarts from 0.
let second_keys: Vec<i64> = vec![100, 200, 300];
let col2 = Arc::new(Int64Array::from(second_keys)) as ArrayRef;
gv.intern(&[col2], &mut groups).unwrap();
assert_eq!(groups, vec![0, 1, 2]);

// size() must be consistent with the live map allocation.
let reported = gv.size();
let map_alloc = gv.map.allocation_size();
assert!(
reported >= map_alloc,
"size() ({reported}) must cover map.allocation_size() ({map_alloc}) after reuse"
);
}
}
Loading
Loading