diff --git a/datafusion/physical-plan/benches/hash_join_semi_anti.rs b/datafusion/physical-plan/benches/hash_join_semi_anti.rs index 40ba41272ca3b..21a28b17bfdab 100644 --- a/datafusion/physical-plan/benches/hash_join_semi_anti.rs +++ b/datafusion/physical-plan/benches/hash_join_semi_anti.rs @@ -15,7 +15,7 @@ // specific language governing permissions and limitations // under the License. -//! Criterion benchmarks for Hash Join with RightSemi/RightAnti joins with Int32 keys. +//! Criterion benchmarks for hash join build sizing and RightSemi/RightAnti joins. //! //! ## Key Benchmark Axes //! @@ -48,9 +48,9 @@ use criterion::{BenchmarkId, Criterion, criterion_group, criterion_main}; use datafusion_common::{JoinType, NullEquality}; use datafusion_execution::TaskContext; use datafusion_physical_expr::expressions::col; -use datafusion_physical_plan::collect; use datafusion_physical_plan::joins::{HashJoinExec, PartitionMode, utils::JoinOn}; use datafusion_physical_plan::test::TestMemoryExec; +use datafusion_physical_plan::{ExecutionPlan, collect}; use tokio::runtime::Runtime; /// Build RecordBatches with Int32 keys. @@ -92,10 +92,7 @@ fn build_batches( batches } -fn make_exec( - batches: &[RecordBatch], - schema: &SchemaRef, -) -> Arc { +fn make_exec(batches: &[RecordBatch], schema: &SchemaRef) -> Arc { TestMemoryExec::try_new_exec(&[batches.to_vec()], Arc::clone(schema), None).unwrap() } @@ -108,8 +105,8 @@ fn schema() -> SchemaRef { } fn do_hash_join( - left: Arc, - right: Arc, + left: Arc, + right: Arc, join_type: JoinType, rt: &Runtime, ) -> usize { @@ -442,5 +439,89 @@ fn bench_hash_join_semi_anti(c: &mut Criterion) { group.finish(); } -criterion_group!(benches, bench_hash_join_semi_anti); +/// Isolate generic hash-table construction from probe fanout and ArrayMap selection. +fn bench_hash_join_build(c: &mut Criterion) { + const ROWS: usize = 1_000_000; + let rt = Runtime::new().unwrap(); + let schema = Arc::new(Schema::new(vec![Field::new("key", DataType::Utf8, false)])); + let probe = RecordBatch::try_new( + Arc::clone(&schema), + vec![Arc::new(StringArray::from(vec!["absent"]))], + ) + .unwrap(); + let mut group = c.benchmark_group("hash_join_build"); + + for distribution in [ + "duplicates", + "mid_20k", + "mid_100k", + "mostly_unique", + "unique", + "unique_prefix", + "duplicate_prefix", + ] { + // Allocate independent arrays so retained slices do not inflate memory metrics. + let build: Vec<_> = (0..ROWS) + .step_by(8192) + .map(|start| { + let keys = StringArray::from_iter_values( + (start..(start + 8192).min(ROWS)).map(|row| { + let key = match distribution { + "duplicates" => row % 64, + "mid_20k" => row % 20_000, + "mid_100k" => row % 100_000, + "mostly_unique" => row % 900_000, + "unique" => row, + _ => { + // Identical key multisets in opposite input order. + let row = if distribution == "duplicate_prefix" { + ROWS - 1 - row + } else { + row + }; + if row < ROWS / 2 { row + 64 } else { row % 64 } + } + }; + format!("key_{key}") + }), + ); + RecordBatch::try_new(Arc::clone(&schema), vec![Arc::new(keys)]).unwrap() + }) + .collect(); + let run = || { + let join = Arc::new( + HashJoinExec::try_new( + make_exec(&build, &schema), + make_exec(std::slice::from_ref(&probe), &schema), + vec![(col("key", &schema).unwrap(), col("key", &schema).unwrap())], + None, + &JoinType::Inner, + None, + PartitionMode::CollectLeft, + NullEquality::NullEqualsNothing, + false, + ) + .unwrap(), + ); + let output = rt + .block_on(collect(join.clone(), Arc::new(TaskContext::default()))) + .unwrap(); + let metrics = join.metrics().unwrap(); + ( + output.iter().map(RecordBatch::num_rows).sum::(), + metrics.sum_by_name("build_input_rows").unwrap().as_usize(), + metrics.sum_by_name("build_mem_used").unwrap().as_usize(), + ) + }; + // An absent probe avoids materializing duplicate matches, but must build the index. + let (output_rows, build_rows, build_bytes) = run(); + assert_eq!(output_rows, 0); + assert_eq!(build_rows, ROWS); + eprintln!("hash_join_build/{distribution}: build_mem_used={build_bytes} bytes"); + group.bench_function(distribution, |b| b.iter(run)); + } + group.finish(); +} + +criterion_group!(benches, bench_hash_join_semi_anti, bench_hash_join_build); criterion_main!(benches); diff --git a/datafusion/physical-plan/src/joins/hash_join/compact_hash_map.rs b/datafusion/physical-plan/src/joins/hash_join/compact_hash_map.rs new file mode 100644 index 0000000000000..d3eb30666943f --- /dev/null +++ b/datafusion/physical-plan/src/joins/hash_join/compact_hash_map.rs @@ -0,0 +1,425 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Adapt hash lookup capacity to the build keys while retaining every row index. + +use std::fmt; +use std::hash::BuildHasher; +use std::mem::size_of; + +use arrow::array::AsArray; +use arrow::datatypes::DataType; +use arrow::record_batch::RecordBatch; +use datafusion_common::hash_utils::{RandomState, create_hashes}; +use datafusion_common::utils::memory::estimate_memory_size; +use datafusion_common::{NullEquality, Result}; +use datafusion_execution::memory_pool::MemoryReservation; +use datafusion_physical_expr::PhysicalExprRef; +use datafusion_physical_expr::expressions::Column; +use datafusion_physical_expr_common::utils::evaluate_expressions_to_arrays; +use hashbrown::HashTable; + +use crate::joins::join_hash_map::update_from_iter_inner; +use crate::joins::utils::matchable_join_keys; + +/// Bound hash scratch and the number of potentially new hashes per insertion. +const HASH_BUILD_CHUNK_ROWS: usize = 8192; + +// Only a capacity hint: every row still goes through the normal build path. +// Sample directly from flat byte columns, without copying values or evaluating +// expressions again. Other key representations keep the existing growth policy. +fn sampled_capacity( + batches: &[RecordBatch], + on: &[PhysicalExprRef], + num_rows: usize, + random_state: &RandomState, + reservation: &MemoryReservation, +) -> Option<(usize, usize)> { + const SAMPLES: usize = 2048; + if num_rows < 128 * SAMPLES || on.len() != 1 { + return None; + } + let column = on[0].downcast_ref::()?; + let data_type = batches.first()?.column(column.index()).data_type(); + if !matches!( + data_type, + DataType::Utf8 + | DataType::LargeUtf8 + | DataType::Utf8View + | DataType::Binary + | DataType::LargeBinary + | DataType::BinaryView + ) { + return None; + } + // Sparse non-null samples cannot estimate the eligible population safely. + // Abstain before allocating scratch so the original path and peak are kept. + if batches + .iter() + .any(|batch| batch.column(column.index()).null_count() != 0) + { + return None; + } + let sample_reservation = reservation.new_empty(); + let sample_bytes = SAMPLES * size_of::(); + sample_reservation.try_grow(sample_bytes).ok()?; + let mut samples = Vec::with_capacity(SAMPLES); + for i in 0..SAMPLES { + samples.push(random_state.hash_one(i) % num_rows as u64); + } + // Sample across the whole input, including within grouped runs of equal + // values. Repeated row positions must not create false duplicate hashes. + samples.sort_unstable(); + samples.dedup(); + let mut batch_index = 0; + let mut batch_start = 0; + for sample in &mut samples { + let row = *sample as usize; + while row >= batch_start + batches[batch_index].num_rows() { + batch_start += batches[batch_index].num_rows(); + batch_index += 1; + } + let array = batches[batch_index].column(column.index()); + let row = row - batch_start; + let hash = match data_type { + DataType::Utf8 => random_state.hash_one(array.as_string::().value(row)), + DataType::LargeUtf8 => { + random_state.hash_one(array.as_string::().value(row)) + } + DataType::Utf8View => { + random_state.hash_one(array.as_string_view().value(row)) + } + DataType::Binary => { + random_state.hash_one(array.as_binary::().value(row)) + } + DataType::LargeBinary => { + random_state.hash_one(array.as_binary::().value(row)) + } + DataType::BinaryView => { + random_state.hash_one(array.as_binary_view().value(row)) + } + _ => unreachable!(), + }; + *sample = hash; + } + samples.sort_unstable(); + let mut pairs = 0_u64; + let mut maximum_frequency = 0; + let mut singletons = 0; + let mut distinct = 0; + let mut triples = 0_u64; + let mut start = 0; + while start < samples.len() { + let mut end = start + 1; + while end < samples.len() && samples[start] == samples[end] { + end += 1; + } + let count = end - start; + maximum_frequency = maximum_frequency.max(count); + let count = count as u64; + pairs += count * (count - 1) / 2; + triples += count * count.saturating_sub(1) * count.saturating_sub(2) / 6; + singletons += usize::from(count == 1); + distinct += 1; + start = end; + } + // Few repeated pairs give an unstable estimate. A frequent sampled key + // indicates skew, where collision-based estimates can badly undercount the + // long tail of unique keys. Abstain for skew, retaining the initial small + // table and its usual full preallocation on the first observed growth. + let entries = if singletons <= samples.len() / 32 { + // Nearly all sampled rows repeat. Retain the bounded starting table + // and let observed growth catch any unsampled tail. + distinct * 2 + } else if maximum_frequency > 8 + || (triples >= 8 && triples * samples.len() as u64 > pairs * pairs) + { + 0 + } else if pairs < 8 { + num_rows + } else { + // Inflate the uniform-key collision estimate to leave growth headroom. + let sample_count = samples.len() as u64; + (sample_count * sample_count.saturating_sub(1) / pairs).min(num_rows as u64) + as usize + }; + Some((entries, sample_bytes)) +} + +// Only flat column arrays can be sliced without repeatedly hashing unused +// dictionary values or changing the input batch seen by computed expressions. +fn is_flat_join_type(data_type: &DataType) -> bool { + data_type.is_primitive() + || matches!( + data_type, + DataType::Null + | DataType::Boolean + | DataType::FixedSizeBinary(_) + | DataType::Utf8 + | DataType::LargeUtf8 + | DataType::Binary + | DataType::LargeBinary + | DataType::Utf8View + | DataType::BinaryView + ) +} + +// Start with one chunk so low-cardinality builds stay small. At the first +// growth, a bounded byte-column sample can replace row-count preallocation. +// If a sample is unavailable or cannot fit with its compaction workspace, try +// the row-count bound, then grow by observed hashes plus the next chunk. +#[expect(clippy::type_complexity)] +pub(super) fn build_compact_hash_map( + batches: &[RecordBatch], + on: &[PhysicalExprRef], + num_rows: usize, + random_state: &RandomState, + null_equality: NullEquality, + reservation: &MemoryReservation, + peak: &mut usize, +) -> Result<(HashTable<(u64, T)>, Vec)> +where + T: Copy + Default + TryFrom + PartialOrd, + >::Error: fmt::Debug, +{ + let initial_reserved = reservation.size(); + let fixed_bytes = size_of::>() + size_of::>(); + let chain_bytes = num_rows.checked_mul(size_of::()).ok_or_else(|| { + datafusion_common::exec_datafusion_err!("Hash join row-index size overflow") + })?; + reservation.try_grow(fixed_bytes + chain_bytes)?; + *peak = (*peak).max(reservation.size() - initial_reserved); + let mut next = vec![T::default(); num_rows]; + let mut table = HashTable::new(); + let scratch = reservation.new_empty(); + let max_batch_rows = batches.iter().map(RecordBatch::num_rows).max().unwrap_or(0); + let chunk_rows = max_batch_rows.min(HASH_BUILD_CHUNK_ROWS); + let mut hash_chunk_rows = chunk_rows; + if let Some(batch) = batches.first() { + for expr in on { + if !expr.is::() + || !is_flat_join_type(&expr.data_type(&batch.schema())?) + { + // Slicing a dictionary retains all its values. Hash complex keys + // and evaluate computed keys once per original batch. + hash_chunk_rows = max_batch_rows; + break; + } + } + } + scratch.try_grow(hash_chunk_rows * size_of::())?; + let mut hashes = vec![0; hash_chunk_rows]; + *peak = (*peak).max(reservation.size() - initial_reserved + scratch.size()); + let mut tried_preallocation = false; + let mut sampled_preallocation = false; + let mut compact_sampled_growth = false; + // Held capacity is part of the reservation peak, without allocating it yet. + let mut compaction_headroom = 0; + 'build: loop { + let mut offset = 0; + for batch in batches.iter().rev() { + for hash_start in (0..batch.num_rows()).step_by(hash_chunk_rows.max(1)).rev() + { + let hash_rows = (batch.num_rows() - hash_start).min(hash_chunk_rows); + let chunk = (hash_rows < batch.num_rows()) + .then(|| batch.slice(hash_start, hash_rows)); + let keys = + evaluate_expressions_to_arrays(on, chunk.as_ref().unwrap_or(batch))?; + hashes[..hash_rows].fill(0); + let hashes = + create_hashes(&keys, random_state, &mut hashes[..hash_rows])?; + let valid = matchable_join_keys(&keys, null_equality); + for start in (0..hash_rows).step_by(HASH_BUILD_CHUNK_ROWS).rev() { + let rows = (hash_rows - start).min(HASH_BUILD_CHUNK_ROWS); + let hashes = &hashes[start..start + rows]; + let valid = valid.as_ref().map(|valid| valid.slice(start, rows)); + let additional = rows - valid.as_ref().map_or(0, |n| n.null_count()); + if additional > table.capacity() - table.len() { + // Keep the old allocation charged during rehash. The fixed-size + // allowance covers control-group padding; at least eight elements + // also covers hashbrown's minimum bucket sizes. + let minimum = (table.len() + additional).max(chunk_rows); + let mut entries = if table.capacity() == 0 { + minimum + } else if !tried_preallocation || sampled_preallocation { + // Sample only when the bounded table first outgrows + // its capacity. A zero hint or unsupported key keeps + // the original row-count preallocation attempt. + reservation.shrink(std::mem::take(&mut compaction_headroom)); + let estimate = if !tried_preallocation { + sampled_capacity( + batches, + on, + num_rows, + random_state, + reservation, + ) + } else { + None + }; + tried_preallocation = true; + sampled_preallocation = false; + if let Some((estimate, sample_bytes)) = estimate { + *peak = (*peak).max( + reservation.size() - initial_reserved + + scratch.size() + + sample_bytes, + ); + if estimate == 0 { + num_rows + } else { + let estimate = estimate + .saturating_add(chunk_rows) + .min(num_rows) + .max(minimum); + sampled_preallocation = estimate < num_rows; + estimate + } + } else { + num_rows + } + } else { + minimum + }; + let mut bytes = estimate_memory_size::<(u64, T)>( + entries.max(8), + fixed_bytes, + )?; + // Hints in the full table's bucket tier need no separate + // compaction workspace or sampled-growth policy. + if sampled_preallocation + && estimate_memory_size::<(u64, T)>( + num_rows.max(8), + fixed_bytes, + ) + .is_ok_and(|full_bytes| bytes == full_bytes) + { + entries = num_rows; + sampled_preallocation = false; + } + let mut headroom = if sampled_preallocation { + // A table at most half full shrinks to at most half its + // buckets. Keep that workspace until the hint is + // outgrown or compacted, so a sample cannot consume + // memory needed to release its excess allocation. + ((bytes - fixed_bytes) / 2 + fixed_bytes) + .max(estimate_memory_size::<(u64, T)>(8, fixed_bytes)?) + } else { + 0 + }; + let mut admission = bytes + .checked_add(headroom) + .ok_or_else(|| { + datafusion_common::exec_datafusion_err!( + "Hash join compaction reservation size overflow" + ) + }) + .and_then(|bytes| reservation.try_grow(bytes)); + if admission.is_err() && headroom != 0 { + // The full table may fit even when the sampled table + // plus guaranteed compaction workspace does not. + headroom = 0; + sampled_preallocation = false; + entries = num_rows; + bytes = estimate_memory_size::<(u64, T)>( + entries.max(8), + fixed_bytes, + )?; + admission = reservation.try_grow(bytes); + } + if admission.is_err() && entries != minimum { + entries = minimum; + bytes = estimate_memory_size::<(u64, T)>( + entries.max(8), + fixed_bytes, + )?; + admission = reservation.try_grow(bytes); + } + compaction_headroom = headroom; + let replay = admission.is_err(); + if let Err(error) = admission { + // If only the replacement fits, drop the partial index + // before releasing its charge, then rebuild from the start. + if table.capacity() == 0 { + return Err(error); + } + let old_bytes = table.allocation_size(); + table = HashTable::new(); + reservation.shrink(old_bytes); + reservation.try_grow(bytes)?; + } + *peak = (*peak) + .max(reservation.size() - initial_reserved + scratch.size()); + let old_bytes = table.allocation_size(); + table + .try_reserve(entries - table.len(), |&(hash, _)| hash) + .map_err(|e| { + datafusion_common::DataFusionError::ResourcesExhausted( + format!("Hash join table allocation: {e}"), + ) + })?; + reservation.shrink(old_bytes + bytes - table.allocation_size()); + if sampled_preallocation { + compact_sampled_growth = true; + } else if entries == num_rows { + compact_sampled_growth = false; + } + if replay { + next.fill(T::default()); + continue 'build; + } + } + let row_offset = offset + hash_start + start; + let iter = hashes + .iter() + .enumerate() + .filter(|(i, _)| valid.as_ref().is_none_or(|n| n.is_valid(*i))) + .map(|(i, hash)| (row_offset + i, hash)); + update_from_iter_inner(&mut table, &mut next, iter.rev(), 0); + } + } + offset += batch.num_rows(); + } + break; + } + drop(hashes); + drop(scratch); + // Keep minimum growth after an outgrown sample compact too. A replacement + // that fit beside the original row-count-sized table also fits beside this + // smaller table. Successful full preallocation keeps its original threshold. + let should_compact = if compact_sampled_growth { + table.capacity() / 2 >= table.len() + } else { + table.capacity() / 4 > table.len() + }; + if should_compact { + let bytes = estimate_memory_size::<(u64, T)>(table.len().max(8), fixed_bytes)?; + let additional = bytes.saturating_sub(compaction_headroom); + if additional == 0 || reservation.try_grow(additional).is_ok() { + *peak = (*peak).max(reservation.size() - initial_reserved); + let old_bytes = table.allocation_size(); + table.shrink_to_fit(|&(hash, _)| hash); + reservation.shrink(old_bytes + bytes - table.allocation_size()); + compaction_headroom = compaction_headroom.saturating_sub(bytes); + } + } + reservation.shrink(compaction_headroom); + Ok((table, next)) +} + +#[cfg(test)] +mod tests; diff --git a/datafusion/physical-plan/src/joins/hash_join/compact_hash_map/tests.rs b/datafusion/physical-plan/src/joins/hash_join/compact_hash_map/tests.rs new file mode 100644 index 0000000000000..782a0181a576c --- /dev/null +++ b/datafusion/physical-plan/src/joins/hash_join/compact_hash_map/tests.rs @@ -0,0 +1,1362 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use super::*; + +use std::sync::Arc; +use std::sync::atomic::{AtomicUsize, Ordering}; + +use arrow::array::{ArrayRef, DictionaryArray, Int32Array, StringArray, StructArray}; +use arrow::compute::concat_batches; +use arrow::datatypes::{Field, Int32Type, Schema}; +use datafusion_common::cast::as_int32_array; +use datafusion_common::utils::memory::RecordBatchMemoryCounter; +use datafusion_common::{DataFusionError, JoinType}; +use datafusion_execution::TaskContext; +use datafusion_execution::memory_pool::{ + GreedyMemoryPool, MemoryConsumer, MemoryPool, PeakRecordingPool, +}; +use datafusion_execution::runtime_env::RuntimeEnvBuilder; +use datafusion_expr::{Volatility, create_udf}; +use datafusion_physical_expr::ScalarFunctionExpr; + +use super::super::exec::HASH_JOIN_SEED; +use crate::ExecutionPlan; +use crate::common; +use crate::joins::PartitionMode; +use crate::joins::hash_join::HashJoinExec; +use crate::joins::join_hash_map::{JoinHashMapType, JoinHashMapU32}; +use crate::joins::utils::update_hash; +use crate::test::TestMemoryExec; + +#[tokio::test] +async fn compact_hash_build_with_duplicates_and_nulls() -> Result<()> { + let rows = 65_536; + let schema = Arc::new(Schema::new(vec![ + Field::new("key", DataType::Utf8, true), + Field::new("value", DataType::Int32, false), + ])); + let build = RecordBatch::try_new( + Arc::clone(&schema), + vec![ + Arc::new(StringArray::from_iter( + (0..rows).map(|i| [Some("GA"), None, Some("CA")][i % 3]), + )), + Arc::new(Int32Array::from_iter_values(0..rows as i32)), + ], + )?; + let probe = RecordBatch::try_new( + Arc::clone(&schema), + vec![ + Arc::new(StringArray::from(vec![Some("GA"), None, Some("CA")])), + Arc::new(Int32Array::from(vec![-1, -2, -3])), + ], + )?; + let small_limit = 2 * 1024 * 1024; + assert!( + estimate_memory_size::<(u64, u32)>(rows, size_of::())? + > small_limit + ); + for null_equality in [ + NullEquality::NullEqualsNothing, + NullEquality::NullEqualsNull, + ] { + let mut outputs = vec![]; + for limit in [16 * 1024 * 1024, small_limit] { + let pool: Arc = Arc::new(GreedyMemoryPool::new(limit)); + let runtime = RuntimeEnvBuilder::new() + .with_memory_pool(Arc::clone(&pool)) + .build_arc()?; + let context = Arc::new(TaskContext::default().with_runtime(runtime)); + let left = TestMemoryExec::try_new_exec( + &[vec![ + build.slice(0, 17_003), + build.slice(17_003, rows - 17_003 - 1), + build.slice(rows - 1, 1), + ]], + Arc::clone(&schema), + None, + )?; + let right = TestMemoryExec::try_new_exec( + &[vec![probe.clone()]], + Arc::clone(&schema), + None, + )?; + let join = HashJoinExec::try_new( + left, + right, + vec![( + Arc::new(Column::new("key", 0)), + Arc::new(Column::new("key", 0)), + )], + None, + &JoinType::Inner, + None, + PartitionMode::Partitioned, + null_equality, + false, + )?; + let batches = common::collect(join.execute(0, context)?).await?; + outputs.push(concat_batches(&join.schema(), &batches)?); + // Even with ample memory, duplicate keys must not allocate a + // row-count-sized bucket array. Uneven batches exercise offsets. + let peak = join + .metrics() + .unwrap() + .sum_by_name("build_mem_used") + .unwrap() + .as_usize(); + assert!(peak < small_limit); + drop(join); + assert_eq!(pool.reserved(), 0); + } + assert_eq!(outputs[0], outputs[1]); + let mut values = as_int32_array(outputs[1].column(1).as_ref())? + .values() + .to_vec(); + values.sort_unstable(); + assert_eq!( + values, + (0..rows as i32) + .filter(|i| null_equality == NullEquality::NullEqualsNull || i % 3 != 1) + .collect::>() + ); + } + Ok(()) +} + +#[tokio::test] +async fn compact_hash_build_leaves_room_for_visited_bitmap() -> Result<()> { + // Each budget admits a different initial bucket allocation, exercising + // compaction when a sampled capacity also needs to release probe headroom. + let mut errors = Vec::new(); + for (distinct_keys, table_rows) in [ + (64, HASH_BUILD_CHUNK_ROWS), + (6144, 2 * HASH_BUILD_CHUNK_ROWS), + (20_000, 4 * HASH_BUILD_CHUNK_ROWS), + ] { + let rows = 1_000_000; + let schema = + Arc::new(Schema::new(vec![Field::new("key", DataType::Utf8, false)])); + let build = (0..rows) + .step_by(HASH_BUILD_CHUNK_ROWS) + .map(|start| { + RecordBatch::try_new( + Arc::clone(&schema), + vec![Arc::new(StringArray::from_iter_values( + (start..(start + HASH_BUILD_CHUNK_ROWS).min(rows)) + .map(|row| format!("key_{}", row % distinct_keys)), + ))], + ) + }) + .collect::, _>>()?; + let mut counter = RecordBatchMemoryCounter::new(); + let input_bytes = build + .iter() + .map(|batch| counter.count_batch(batch)) + .sum::(); + // The table allowance and hash scratch fit, but retaining the entire + // allowance leaves too little room for the one-bit-per-row visited bitmap. + let fixed_bytes = size_of::(); + let limit = input_bytes + + fixed_bytes + + rows * size_of::() + + HASH_BUILD_CHUNK_ROWS * size_of::() + + estimate_memory_size::<(u64, u32)>(table_rows, fixed_bytes)?; + let pool: Arc = Arc::new(GreedyMemoryPool::new(limit)); + let runtime = RuntimeEnvBuilder::new() + .with_memory_pool(Arc::clone(&pool)) + .build_arc()?; + let context = Arc::new(TaskContext::default().with_runtime(runtime)); + let left = TestMemoryExec::try_new_exec(&[build], Arc::clone(&schema), None)?; + let probe = RecordBatch::try_new( + Arc::clone(&schema), + vec![Arc::new(StringArray::from(vec!["absent"]))], + )?; + let right = + TestMemoryExec::try_new_exec(&[vec![probe]], Arc::clone(&schema), None)?; + let join = HashJoinExec::try_new( + left, + right, + vec![( + Arc::new(Column::new("key", 0)), + Arc::new(Column::new("key", 0)), + )], + None, + &JoinType::LeftAnti, + None, + PartitionMode::CollectLeft, + NullEquality::NullEqualsNothing, + false, + )?; + let result = common::collect(join.execute(0, context)?).await; + drop(join); + assert_eq!(pool.reserved(), 0); + let output = match result { + Ok(output) => output, + Err(error) => { + errors.push(format!( + "distinct_keys={distinct_keys}, limit={limit}: {error}" + )); + continue; + } + }; + assert_eq!( + output.iter().map(RecordBatch::num_rows).sum::(), + rows + ); + } + assert!(errors.is_empty(), "{}", errors.join("\n")); + Ok(()) +} + +#[test] +fn compact_hash_build_evaluates_complex_keys_once() -> Result<()> { + let rows = 2 * HASH_BUILD_CHUNK_ROWS + 17; + let dictionary: ArrayRef = Arc::new(DictionaryArray::::try_new( + Int32Array::from_iter( + (0..rows + 3).map(|i| (i % 7 != 0).then_some((i % 64) as i32)), + ), + Arc::new(StringArray::from_iter_values( + (0..rows).map(|i| format!("key-{i}")), + )), + )?); + let dictionary = dictionary.slice(3, rows); + let nested: ArrayRef = Arc::new(StructArray::from(vec![( + Arc::new(Field::new("value", dictionary.data_type().clone(), true)), + Arc::clone(&dictionary), + )])); + for key in [dictionary, nested] { + let batch = RecordBatch::try_from_iter([("key", Arc::clone(&key))])?; + for computed in [false, true] { + let evaluations = Arc::new(AtomicUsize::new(0)); + let counter = Arc::clone(&evaluations); + let udf = create_udf( + "identity", + vec![key.data_type().clone()], + key.data_type().clone(), + Volatility::Immutable, + Arc::new(move |args| { + counter.fetch_add(1, Ordering::Relaxed); + Ok(args[0].clone()) + }), + ); + let expression: PhysicalExprRef = if computed { + Arc::new(ScalarFunctionExpr::new( + "identity", + Arc::new(udf), + vec![Arc::new(Column::new("key", 0))], + Arc::new(Field::new("identity", key.data_type().clone(), true)), + Arc::default(), + )) + } else { + Arc::new(Column::new("key", 0)) + }; + let on = vec![expression]; + let pool: Arc = + Arc::new(GreedyMemoryPool::new(16 * 1024 * 1024)); + let reservation = MemoryConsumer::new("complex key test").register(&pool); + let (table, next) = build_compact_hash_map::( + std::slice::from_ref(&batch), + &on, + rows, + HASH_JOIN_SEED.random_state(), + NullEquality::NullEqualsNothing, + &reservation, + &mut 0, + )?; + assert_eq!(evaluations.load(Ordering::Relaxed), usize::from(computed)); + + // Compare with the existing whole-batch insertion path, including + // rows excluded by dictionary NULLs and nested key hashing. + let compact = JoinHashMapU32::new(table, next); + let mut expected = JoinHashMapU32::with_capacity(rows); + let mut hashes = vec![0; rows]; + update_hash( + &on, + &batch, + &mut expected, + 0, + HASH_JOIN_SEED.random_state(), + &mut hashes, + 0, + true, + NullEquality::NullEqualsNothing, + )?; + // Sampling the probe hashes avoids quadratic output with forced + // hash collisions while still traversing each sampled full chain. + let probe_hashes = hashes + .iter() + .step_by(257) + .take(8) + .copied() + .collect::>(); + assert_eq!( + compact + .get_matched_indices(Box::new(probe_hashes.iter().enumerate()), None), + expected + .get_matched_indices(Box::new(probe_hashes.iter().enumerate()), None), + ); + drop((compact, reservation)); + assert_eq!(pool.reserved(), 0); + } + } + Ok(()) +} + +#[test] +fn compact_hash_build_preserves_fifo_across_batches() -> Result<()> { + let rows = 3 * HASH_BUILD_CHUNK_ROWS + 17; + let batch = RecordBatch::try_from_iter([( + "key", + Arc::new(Int32Array::from_iter( + (0..rows).map(|i| (i % 7 != 0).then_some((i % 37) as i32)), + )) as ArrayRef, + )])?; + let batches = vec![ + batch.slice(0, HASH_BUILD_CHUNK_ROWS + 5), + batch.slice(HASH_BUILD_CHUNK_ROWS + 5, 2 * HASH_BUILD_CHUNK_ROWS + 11), + batch.slice(rows - 1, 1), + ]; + let on = vec![Arc::new(Column::new("key", 0)) as PhysicalExprRef]; + for null_equality in [ + NullEquality::NullEqualsNothing, + NullEquality::NullEqualsNull, + ] { + let pool: Arc = Arc::new(GreedyMemoryPool::new(16 * 1024 * 1024)); + let reservation = MemoryConsumer::new("FIFO compact build").register(&pool); + let (table, next) = build_compact_hash_map::( + &batches, + &on, + rows, + HASH_JOIN_SEED.random_state(), + null_equality, + &reservation, + &mut 0, + )?; + let compact = JoinHashMapU32::new(table, next); + let mut expected = JoinHashMapU32::with_capacity(rows); + let mut offset = 0; + for batch in batches.iter().rev() { + let mut hashes = vec![0; batch.num_rows()]; + update_hash( + &on, + batch, + &mut expected, + offset, + HASH_JOIN_SEED.random_state(), + &mut hashes, + 0, + true, + null_equality, + )?; + offset += batch.num_rows(); + } + let probe = batch.slice(0, 38); + let mut probe_hashes = vec![0; probe.num_rows()]; + create_hashes( + probe.columns(), + HASH_JOIN_SEED.random_state(), + &mut probe_hashes, + )?; + assert_eq!( + compact.get_matched_indices(Box::new(probe_hashes.iter().enumerate()), None), + expected.get_matched_indices(Box::new(probe_hashes.iter().enumerate()), None), + ); + drop((compact, reservation)); + assert_eq!(pool.reserved(), 0); + } + Ok(()) +} + +#[test] +fn compact_hash_build_growth_and_chain_accounting() -> Result<()> { + let rows = 8 * HASH_BUILD_CHUNK_ROWS; + for (rows, distinct_keys) in [ + (rows, rows), + (rows, HASH_BUILD_CHUNK_ROWS), + (rows, 2 * HASH_BUILD_CHUNK_ROWS), + (1_000_000, 100_000), + ] { + let batch = RecordBatch::try_from_iter([( + "key", + Arc::new(Int32Array::from_iter_values( + (0..rows).map(|i| (i % distinct_keys) as i32), + )) as ArrayRef, + )])?; + let on = vec![Arc::new(Column::new("key", 0)) as PhysicalExprRef]; + let fixed_and_chain = size_of::() + rows * size_of::(); + // The replacement fits, but old and new tables cannot coexist. + // Both budgets must produce identical index heads and duplicate chains. + let mut compact_limit = fixed_and_chain + + HASH_BUILD_CHUNK_ROWS * size_of::() + + estimate_memory_size::<(u64, u32)>( + (distinct_keys + HASH_BUILD_CHUNK_ROWS).min(rows), + size_of::(), + )?; + let deny_final_compaction = rows == 1_000_000; + if deny_final_compaction { + // Admit the row-count preallocation alongside the first small table, + // but leave too little headroom for the final compact replacement. + compact_limit = fixed_and_chain + + HASH_BUILD_CHUNK_ROWS * size_of::() + + estimate_memory_size::<(u64, u32)>( + HASH_BUILD_CHUNK_ROWS, + size_of::(), + )? + + estimate_memory_size::<(u64, u32)>(rows, size_of::())?; + } + let mut expected = None; + for limit in [128 * 1024 * 1024, compact_limit] { + let recording = Arc::new(PeakRecordingPool::new(Arc::new( + GreedyMemoryPool::new(limit), + ))); + let pool: Arc = Arc::clone(&recording) as _; + let reservation = MemoryConsumer::new("compact growth test").register(&pool); + let mut peak = 0; + let (table, next) = build_compact_hash_map::( + std::slice::from_ref(&batch), + &on, + rows, + HASH_JOIN_SEED.random_state(), + NullEquality::NullEqualsNothing, + &reservation, + &mut peak, + )?; + assert_eq!( + reservation.size(), + table.allocation_size() + fixed_and_chain + ); + assert!(peak > reservation.size() && peak <= limit); + assert_eq!(peak, recording.peak_reserved()); + // With forced hash collisions, no row-count preallocation occurs. + if deny_final_compaction && table.len() > HASH_BUILD_CHUNK_ROWS { + assert_eq!(table.capacity() / 4 > table.len(), limit == compact_limit); + } + let mut entries = table.iter().copied().collect::>(); + entries.sort_unstable(); + if let Some(expected) = &expected { + assert_eq!(&(entries, next), expected); + } else { + expected = Some((entries, next)); + } + drop((table, reservation)); + assert_eq!(pool.reserved(), 0); + } + } + Ok(()) +} + +#[test] +fn compact_hash_build_compacts_speculative_capacity() -> Result<()> { + fn check() -> Result<()> + where + T: Copy + Default + TryFrom + PartialOrd + Into, + >::Error: fmt::Debug, + { + let rows = 1_000_000; + let distinct_keys = 10_000; + let batch = RecordBatch::try_from_iter([( + "key", + Arc::new(Int32Array::from_iter_values( + (0..rows).map(|i| (i % distinct_keys) as i32), + )) as ArrayRef, + )])?; + let on = vec![Arc::new(Column::new("key", 0)) as PhysicalExprRef]; + let pool: Arc = + Arc::new(GreedyMemoryPool::new(128 * 1024 * 1024)); + let reservation = + MemoryConsumer::new("compact speculative capacity").register(&pool); + let mut peak = 0; + let (table, next) = build_compact_hash_map::( + std::slice::from_ref(&batch), + &on, + rows, + HASH_JOIN_SEED.random_state(), + NullEquality::NullEqualsNothing, + &reservation, + &mut peak, + )?; + let fixed_bytes = size_of::>() + size_of::>(); + // Crossing the first chunk's capacity can speculate on all build rows. + // The retained table should instead reflect the distinct build hashes. + assert!( + table.allocation_size() + <= estimate_memory_size::<(u64, T)>(4 * distinct_keys, fixed_bytes)? + ); + assert_eq!( + reservation.size(), + fixed_bytes + rows * size_of::() + table.allocation_size() + ); + assert_eq!(pool.reserved(), reservation.size()); + assert!(peak > reservation.size()); + + let mut hashes = vec![0; distinct_keys]; + create_hashes( + batch.slice(0, distinct_keys).columns(), + HASH_JOIN_SEED.random_state(), + &mut hashes, + )?; + // Traverse each chain once, including when hashes deliberately collide. + // This checks that compaction preserves every row and its FIFO order. + let mut seen = vec![false; rows]; + for &(hash, head) in table.iter() { + let mut row = head.into() as usize; + while row != 0 { + let index = row - 1; + assert!(!seen[index]); + seen[index] = true; + assert_eq!(hash, hashes[index % distinct_keys]); + let next_row = next[index].into() as usize; + assert!(next_row == 0 || next_row > row); + row = next_row; + } + } + assert!(seen.into_iter().all(|visited| visited)); + drop((table, reservation)); + assert_eq!(pool.reserved(), 0); + Ok(()) + } + + check::()?; + check::() +} + +#[test] +fn compact_hash_build_releases_failed_admissions() -> Result<()> { + let rows = HASH_BUILD_CHUNK_ROWS; + let batch = RecordBatch::try_from_iter([( + "key", + Arc::new(Int32Array::from_iter_values(0..rows as i32)) as ArrayRef, + )])?; + let on = vec![Arc::new(Column::new("key", 0)) as PhysicalExprRef]; + let fixed_and_chain = size_of::() + rows * size_of::(); + // Reject the chain, then scratch, then the first bucket allocation. + for limit in [ + fixed_and_chain - 1, + fixed_and_chain, + fixed_and_chain + rows * size_of::(), + ] { + let pool: Arc = Arc::new(GreedyMemoryPool::new(limit)); + let reservation = MemoryConsumer::new("failed compact build").register(&pool); + let result = build_compact_hash_map::( + std::slice::from_ref(&batch), + &on, + rows, + HASH_JOIN_SEED.random_state(), + NullEquality::NullEqualsNothing, + &reservation, + &mut 0, + ); + assert!(matches!( + result, + Err(DataFusionError::ResourcesExhausted(_)) + )); + drop(reservation); + assert_eq!(pool.reserved(), 0); + } + Ok(()) +} + +#[test] +fn compact_hash_build_index_widths_and_empty_input() -> Result<()> { + let rows = 2 * HASH_BUILD_CHUNK_ROWS + 17; + for count in [0, rows] { + for all_null in [false, true] { + let batch = RecordBatch::try_from_iter([( + "key", + Arc::new(Int32Array::from_iter( + (0..count) + .map(|i| (!all_null && i % 7 != 0).then_some((i % 17) as i32)), + )) as ArrayRef, + )])?; + let on = vec![Arc::new(Column::new("key", 0)) as PhysicalExprRef]; + for null_equality in [ + NullEquality::NullEqualsNothing, + NullEquality::NullEqualsNull, + ] { + let pool: Arc = + Arc::new(GreedyMemoryPool::new(16 * 1024 * 1024)); + let narrow = MemoryConsumer::new("u32 compact build").register(&pool); + let wide = MemoryConsumer::new("u64 compact build").register(&pool); + let (table32, next32) = build_compact_hash_map::( + std::slice::from_ref(&batch), + &on, + count, + HASH_JOIN_SEED.random_state(), + null_equality, + &narrow, + &mut 0, + )?; + let (table64, next64) = build_compact_hash_map::( + std::slice::from_ref(&batch), + &on, + count, + HASH_JOIN_SEED.random_state(), + null_equality, + &wide, + &mut 0, + )?; + let mut entries32 = table32 + .iter() + .map(|&(hash, row)| (hash, u64::from(row))) + .collect::>(); + let mut entries64 = table64.iter().copied().collect::>(); + entries32.sort_unstable(); + entries64.sort_unstable(); + assert_eq!(entries32, entries64); + assert_eq!( + next32.into_iter().map(u64::from).collect::>(), + next64 + ); + if count == 0 + || (all_null && null_equality == NullEquality::NullEqualsNothing) + { + assert!(table32.is_empty()); + } + drop((table32, table64, narrow, wide)); + assert_eq!(pool.reserved(), 0); + } + } + } + Ok(()) +} + +fn sampled_string_batch( + rows: usize, + key: impl Fn(usize) -> Option, +) -> Result { + Ok(RecordBatch::try_from_iter([( + "key", + Arc::new(StringArray::from_iter((0..rows).map(key))) as ArrayRef, + )])?) +} + +// Validate every row once, including when all hashes deliberately collide. +// Returns distinct hashes, retained table allocation, and tracked build peak. +fn assert_sampled_build_preserves_rows( + batches: &[RecordBatch], + on: &[PhysicalExprRef], + null_equality: NullEquality, +) -> Result<(usize, usize, usize)> { + assert_sampled_build_preserves_rows_with_limit( + batches, + on, + null_equality, + 128 * 1024 * 1024, + ) +} + +fn assert_sampled_build_preserves_rows_with_limit( + batches: &[RecordBatch], + on: &[PhysicalExprRef], + null_equality: NullEquality, + limit: usize, +) -> Result<(usize, usize, usize)> { + let rows = batches.iter().map(RecordBatch::num_rows).sum(); + let recording = Arc::new(PeakRecordingPool::new(Arc::new(GreedyMemoryPool::new( + limit, + )))); + let pool: Arc = Arc::clone(&recording) as _; + let reservation = MemoryConsumer::new("sampled string build").register(&pool); + let mut peak = 0; + let (table, next) = build_compact_hash_map::( + batches, + on, + rows, + HASH_JOIN_SEED.random_state(), + null_equality, + &reservation, + &mut peak, + )?; + assert_eq!(peak, recording.peak_reserved()); + assert_eq!( + reservation.size(), + size_of::() + rows * size_of::() + table.allocation_size() + ); + assert_eq!(pool.reserved(), reservation.size()); + let summary = (table.len(), table.allocation_size(), peak); + + // The builder indexes batches in reverse order. Identity expressions in the + // computed-key test produce the same hashes as the underlying string column. + let mut expected_hashes = vec![0; rows]; + let mut expected_valid = vec![false; rows]; + let mut input_order = vec![0; rows]; + let mut offset = 0; + for batch in batches.iter().rev() { + let end = offset + batch.num_rows(); + create_hashes( + batch.columns(), + HASH_JOIN_SEED.random_state(), + &mut expected_hashes[offset..end], + )?; + for row in 0..batch.num_rows() { + expected_valid[offset + row] = null_equality == NullEquality::NullEqualsNull + || !batch.column(0).is_null(row); + input_order[offset + row] = rows - end + row; + } + offset = end; + } + let mut seen = vec![false; rows]; + for &(hash, head) in table.iter() { + let mut row = head as usize; + while row != 0 { + let index = row - 1; + assert!(!seen[index], "duplicate row index {index}"); + seen[index] = true; + assert_eq!(hash, expected_hashes[index]); + let next_row = next[index] as usize; + assert!(next_row == 0 || input_order[next_row - 1] > input_order[index]); + row = next_row; + } + } + for (row, (seen, valid)) in seen.into_iter().zip(expected_valid).enumerate() { + assert_eq!(seen, valid, "row index {row}"); + } + drop((table, next, reservation)); + assert_eq!(pool.reserved(), 0); + Ok(summary) +} + +#[test] +fn sampled_hash_build_limits_uniform_key_memory() -> Result<()> { + let rows = 1_000_000; + let on = vec![Arc::new(Column::new("key", 0)) as PhysicalExprRef]; + for distinct in [20_000, 100_000] { + for grouped in [false, true] { + let batch = sampled_string_batch(rows, |row| { + let key = if grouped { + row / (rows / distinct) + } else { + row % distinct + }; + Some(format!("key_{key}")) + })?; + let (_, retained, peak) = assert_sampled_build_preserves_rows( + &[batch], + &on, + NullEquality::NullEqualsNothing, + )?; + let row_sized = + estimate_memory_size::<(u64, u32)>(rows, size_of::())?; + // The entire construction should fit below the old bucket allocation + // alone; final shrinking must not merely hide an oversized build peak. + assert!(peak < row_sized, "distinct={distinct}, grouped={grouped}"); + assert!(retained < row_sized / 4); + assert!( + retained + <= estimate_memory_size::<(u64, u32)>( + 4 * distinct, + size_of::(), + )? + ); + } + } + Ok(()) +} + +#[test] +fn sampled_hash_build_preserves_unique_prefix() -> Result<()> { + let rows = 300_000; + let on = vec![Arc::new(Column::new("key", 0)) as PhysicalExprRef]; + for reverse in [false, true] { + let batch = sampled_string_batch(rows, |row| { + let row = if reverse { rows - 1 - row } else { row }; + let key = if row < rows / 2 { row + 64 } else { row % 64 }; + Some(format!("key_{key}")) + })?; + let pool: Arc = Arc::new(GreedyMemoryPool::new(1024 * 1024)); + let reservation = MemoryConsumer::new("skewed sample").register(&pool); + let (estimate, _) = sampled_capacity( + std::slice::from_ref(&batch), + &on, + rows, + HASH_JOIN_SEED.random_state(), + &reservation, + ) + .expect("large flat string keys can be sampled"); + // Strong skew keeps the initial allocation small, then actual distinct + // hashes trigger row-count preallocation for the long unique tail. + assert!(estimate <= HASH_BUILD_CHUNK_ROWS); + assert_eq!(pool.reserved(), 0); + let (distinct, _, peak) = assert_sampled_build_preserves_rows( + &[batch], + &on, + NullEquality::NullEqualsNothing, + )?; + if distinct > 1 { + assert!( + peak >= estimate_memory_size::<(u64, u32)>( + rows, + size_of::() + )? + ); + } + } + Ok(()) +} + +#[test] +fn sampled_hash_build_preserves_skewed_unique_tail() -> Result<()> { + let rows = 1_000_000; + let unique_rows = 25_000; + let final_duplicate_rows = HASH_BUILD_CHUNK_ROWS; + let on = vec![Arc::new(Column::new("key", 0)) as PhysicalExprRef]; + let batch = sampled_string_batch(rows, |row| { + // The builder visits the unique region near the end. One final chunk + // of repeated keys then exceeds the sampled table's spare capacity. + let key = if (final_duplicate_rows..final_duplicate_rows + unique_rows) + .contains(&row) + { + row - final_duplicate_rows + 1 + } else { + 0 + }; + Some(format!("key_{key}")) + })?; + let pool: Arc = Arc::new(GreedyMemoryPool::new(1024 * 1024)); + let reservation = MemoryConsumer::new("underestimated unique tail").register(&pool); + let (estimate, _) = sampled_capacity( + std::slice::from_ref(&batch), + &on, + rows, + HASH_JOIN_SEED.random_state(), + &reservation, + ) + .expect("large flat string keys can be sampled"); + assert!(estimate > 0 && estimate < unique_rows, "hint={estimate}"); + assert_eq!(pool.reserved(), 0); + let (distinct, _, peak) = assert_sampled_build_preserves_rows( + std::slice::from_ref(&batch), + &on, + NullEquality::NullEqualsNothing, + )?; + // With forced hash collisions the actual index never outgrows its hint. + // Otherwise the builder falls back to row-count preallocation before final + // compaction, preserving all the rows checked by the helper above. + if distinct > 1 { + assert_eq!(distinct, unique_rows + 1); + assert!( + peak >= estimate_memory_size::<(u64, u32)>( + rows, + size_of::(), + )? + ); + } + // The final table fits, but cannot coexist with the initial table during + // growth. Recovering from an underestimated hint must also work on replay. + let fixed_bytes = size_of::(); + let limit = fixed_bytes + + rows * size_of::() + + HASH_BUILD_CHUNK_ROWS * size_of::() + + estimate_memory_size::<(u64, u32)>( + unique_rows + HASH_BUILD_CHUNK_ROWS, + fixed_bytes, + )?; + assert_sampled_build_preserves_rows_with_limit( + &[batch], + &on, + NullEquality::NullEqualsNothing, + limit, + )?; + Ok(()) +} + +#[test] +fn sampled_hash_build_preserves_broad_skew() -> Result<()> { + let rows = 300_000; + let on = vec![Arc::new(Column::new("key", 0)) as PhysicalExprRef]; + for reverse in [false, true] { + let batch = sampled_string_batch(rows, |row| { + let row = if reverse { rows - 1 - row } else { row }; + // Repeated values spread across a broad domain need not look like + // heavy hitters. A mistaken estimate must preserve the unique half. + let key = if row < rows / 2 { + row + 8000 + } else { + row % 8000 + }; + Some(format!("key_{key}")) + })?; + assert_sampled_build_preserves_rows( + &[batch], + &on, + NullEquality::NullEqualsNothing, + )?; + } + Ok(()) +} + +#[test] +fn sampled_hash_build_handles_nulls_and_empty_batches() -> Result<()> { + let rows = 1_000_000; + let on = vec![Arc::new(Column::new("key", 0)) as PhysicalExprRef]; + let batch = sampled_string_batch(rows, |row| { + (row % 1000 == 0).then(|| format!("key_{}", row / 1000)) + })?; + let split = rows / 3 + 1; + let batches = [ + batch.slice(0, 0), + batch.slice(0, split), + batch.slice(split, 0), + batch.slice(split, rows - split), + batch.slice(rows, 0), + ]; + let recording = Arc::new(PeakRecordingPool::new(Arc::new(GreedyMemoryPool::new( + 1024 * 1024, + )))); + let pool: Arc = Arc::clone(&recording) as _; + let reservation = MemoryConsumer::new("nullable sample").register(&pool); + assert!( + sampled_capacity( + &batches, + &on, + rows, + HASH_JOIN_SEED.random_state(), + &reservation, + ) + .is_none() + ); + // Abstain before allocating sample scratch: sampling mostly NULL rows must + // not turn a small valid-key domain into row-count preallocation. + assert_eq!(recording.peak_reserved(), 0); + assert_eq!(pool.reserved(), 0); + let fixed_bytes = size_of::(); + let initial_peak = fixed_bytes + + rows * size_of::() + + HASH_BUILD_CHUNK_ROWS * size_of::() + + estimate_memory_size::<(u64, u32)>(HASH_BUILD_CHUNK_ROWS, fixed_bytes)?; + for equality in [ + NullEquality::NullEqualsNothing, + NullEquality::NullEqualsNull, + ] { + let (_, _, peak) = assert_sampled_build_preserves_rows_with_limit( + &batches, + &on, + equality, + initial_peak, + )?; + assert_eq!(peak, initial_peak); + } + Ok(()) +} + +#[test] +fn sampled_hash_build_handles_flat_byte_representations() -> Result<()> { + use arrow::array::{ + BinaryArray, BinaryViewArray, LargeBinaryArray, LargeStringArray, StringViewArray, + }; + + let rows = 300_000; + let batch = sampled_string_batch(rows, |row| Some(format!("key_{}", row % 20_000)))?; + let source = batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + let arrays: [ArrayRef; 6] = [ + Arc::clone(batch.column(0)), + Arc::new(LargeStringArray::from_iter(source.iter())), + Arc::new(StringViewArray::from_iter(source.iter())), + Arc::new(BinaryArray::from_iter( + source.iter().map(|value| value.map(str::as_bytes)), + )), + Arc::new(LargeBinaryArray::from_iter( + source.iter().map(|value| value.map(str::as_bytes)), + )), + Arc::new(BinaryViewArray::from_iter( + source.iter().map(|value| value.map(str::as_bytes)), + )), + ]; + let on = vec![Arc::new(Column::new("key", 0)) as PhysicalExprRef]; + for array in arrays { + // A nullable schema with no actual NULLs remains eligible. Sampling must + // honor array offsets and empty batches for each flat representation. + let schema = Arc::new(Schema::new(vec![Field::new( + "key", + array.data_type().clone(), + true, + )])); + let batch = RecordBatch::try_new(schema, vec![array.slice(3, rows - 7)])?; + let rows = batch.num_rows(); + let split = rows / 3 + 1; + let batches = [ + batch.slice(0, 0), + batch.slice(0, split), + batch.slice(split, 0), + batch.slice(split, rows - split), + batch.slice(rows, 0), + ]; + let pool: Arc = Arc::new(GreedyMemoryPool::new(1024 * 1024)); + let reservation = MemoryConsumer::new("flat byte sample").register(&pool); + assert!( + sampled_capacity( + &batches, + &on, + rows, + HASH_JOIN_SEED.random_state(), + &reservation, + ) + .is_some() + ); + assert_eq!(pool.reserved(), 0); + assert_sampled_build_preserves_rows( + &batches, + &on, + NullEquality::NullEqualsNothing, + )?; + } + Ok(()) +} + +#[test] +fn sampled_capacity_excludes_dictionary_and_multiple_keys() -> Result<()> { + let rows = 300_000; + let values: ArrayRef = Arc::new(StringArray::from_iter_values( + (0..64).map(|key| format!("key_{key}")), + )); + let dictionary: ArrayRef = Arc::new(DictionaryArray::::try_new( + Int32Array::from_iter_values((0..rows).map(|row| (row % 64) as i32)), + values, + )?); + let dictionary_batch = RecordBatch::try_from_iter([("key", dictionary)])?; + let strings = sampled_string_batch(rows, |row| Some(format!("key_{}", row % 64)))?; + let multi_key_batch = RecordBatch::try_from_iter([ + ("key", Arc::clone(strings.column(0))), + ("other", Arc::clone(strings.column(0))), + ])?; + let pool: Arc = Arc::new(GreedyMemoryPool::new(1024 * 1024)); + let reservation = MemoryConsumer::new("unsupported sample keys").register(&pool); + let column = Arc::new(Column::new("key", 0)) as PhysicalExprRef; + for (batch, on) in [ + (dictionary_batch, vec![Arc::clone(&column)]), + ( + multi_key_batch, + vec![column, Arc::new(Column::new("other", 1)) as PhysicalExprRef], + ), + ] { + assert!( + sampled_capacity( + &[batch], + &on, + rows, + HASH_JOIN_SEED.random_state(), + &reservation, + ) + .is_none() + ); + assert_eq!(pool.reserved(), 0); + } + Ok(()) +} + +#[test] +fn sampled_hash_build_does_not_evaluate_computed_keys_again() -> Result<()> { + let rows = 300_000; + let batch = sampled_string_batch(rows, |row| Some(format!("key_{}", row % 20_000)))?; + let evaluations = Arc::new(AtomicUsize::new(0)); + let counter = Arc::clone(&evaluations); + let udf = create_udf( + "identity", + vec![DataType::Utf8], + DataType::Utf8, + Volatility::Immutable, + Arc::new(move |args| { + counter.fetch_add(1, Ordering::Relaxed); + Ok(args[0].clone()) + }), + ); + let expression: PhysicalExprRef = Arc::new(ScalarFunctionExpr::new( + "identity", + Arc::new(udf), + vec![Arc::new(Column::new("key", 0))], + Arc::new(Field::new("identity", DataType::Utf8, true)), + Arc::default(), + )); + let on = vec![expression]; + let pool: Arc = Arc::new(GreedyMemoryPool::new(1024 * 1024)); + let reservation = MemoryConsumer::new("computed key sample").register(&pool); + assert!( + sampled_capacity( + std::slice::from_ref(&batch), + &on, + rows, + HASH_JOIN_SEED.random_state(), + &reservation, + ) + .is_none() + ); + assert_eq!(evaluations.load(Ordering::Relaxed), 0); + assert_eq!(pool.reserved(), 0); + assert_sampled_build_preserves_rows(&[batch], &on, NullEquality::NullEqualsNothing)?; + assert_eq!(evaluations.load(Ordering::Relaxed), 1); + Ok(()) +} + +#[test] +fn sampled_capacity_releases_failed_admission() -> Result<()> { + let rows = 300_000; + let batch = sampled_string_batch(rows, |row| Some(format!("key_{}", row % 20_000)))?; + let on = vec![Arc::new(Column::new("key", 0)) as PhysicalExprRef]; + let pool: Arc = Arc::new(GreedyMemoryPool::new(1024)); + let reservation = MemoryConsumer::new("exhausted sample pool").register(&pool); + reservation.try_grow(1024)?; + assert!( + sampled_capacity( + std::slice::from_ref(&batch), + &on, + rows, + HASH_JOIN_SEED.random_state(), + &reservation, + ) + .is_none() + ); + assert_eq!(reservation.size(), 1024); + assert_eq!(pool.reserved(), 1024); + drop(reservation); + assert_eq!(pool.reserved(), 0); + + // Exercise cleanup when sampling is denied, and when sampling succeeds but + // neither the hinted table nor its bounded fallback can be admitted. + let scratch_bytes = HASH_BUILD_CHUNK_ROWS * size_of::(); + let fixed_bytes = size_of::(); + let initial_peak = fixed_bytes + + rows * size_of::() + + scratch_bytes + + estimate_memory_size::<(u64, u32)>(HASH_BUILD_CHUNK_ROWS, fixed_bytes)?; + for limit in [initial_peak, initial_peak + scratch_bytes] { + let recording = Arc::new(PeakRecordingPool::new(Arc::new( + GreedyMemoryPool::new(limit), + ))); + let pool: Arc = Arc::clone(&recording) as _; + let reservation = MemoryConsumer::new("failed sampled build").register(&pool); + let mut peak = 0; + let result = build_compact_hash_map::( + std::slice::from_ref(&batch), + &on, + rows, + HASH_JOIN_SEED.random_state(), + NullEquality::NullEqualsNothing, + &reservation, + &mut peak, + ); + match result { + Ok((table, next)) => { + // Forced collisions keep every row in one chain, so the + // initial table never needs growth or sampling. + assert_eq!(table.len(), 1); + drop((table, next)); + } + Err(error) => { + assert!(matches!(error, DataFusionError::ResourcesExhausted(_))); + assert_eq!(peak > initial_peak, limit > initial_peak); + } + } + assert_eq!(peak, recording.peak_reserved()); + drop(reservation); + assert_eq!(pool.reserved(), 0); + } + Ok(()) +} + +#[test] +fn sampled_hash_build_preserves_index_widths() -> Result<()> { + type NormalizedHashIndex = (Vec<(u64, u64)>, Vec); + + fn check( + batch: &RecordBatch, + on: &[PhysicalExprRef], + hashes: &[u64], + distinct_keys: usize, + ) -> Result + where + T: Copy + Default + TryFrom + PartialOrd + Into, + >::Error: fmt::Debug, + { + let rows = batch.num_rows(); + let recording = Arc::new(PeakRecordingPool::new(Arc::new( + GreedyMemoryPool::new(128 * 1024 * 1024), + ))); + let pool: Arc = Arc::clone(&recording) as _; + let reservation = MemoryConsumer::new("sampled index widths").register(&pool); + let mut peak = 0; + let (table, next) = build_compact_hash_map::( + std::slice::from_ref(batch), + on, + rows, + HASH_JOIN_SEED.random_state(), + NullEquality::NullEqualsNothing, + &reservation, + &mut peak, + )?; + let fixed_bytes = size_of::>() + size_of::>(); + assert_eq!(peak, recording.peak_reserved()); + assert_eq!( + reservation.size(), + fixed_bytes + rows * size_of::() + table.allocation_size() + ); + assert_eq!(pool.reserved(), reservation.size()); + if distinct_keys < rows { + assert!(peak < estimate_memory_size::<(u64, T)>(rows, fixed_bytes)?); + } + + let mut seen = vec![false; rows]; + for &(hash, head) in table.iter() { + let mut row = head.into() as usize; + while row != 0 { + let index = row - 1; + assert!(!seen[index]); + seen[index] = true; + assert_eq!(hash, hashes[index]); + let next_row = next[index].into() as usize; + assert!(next_row == 0 || next_row > row); + row = next_row; + } + } + assert!(seen.into_iter().all(|visited| visited)); + let mut entries = table + .iter() + .map(|&(hash, row)| (hash, row.into())) + .collect::>(); + entries.sort_unstable(); + let normalized_next = next.iter().map(|&row| row.into()).collect(); + drop((table, next, reservation)); + assert_eq!(pool.reserved(), 0); + Ok((entries, normalized_next)) + } + + // Both fixtures are large enough for sampling. Repeated keys exercise a + // partial capacity hint; unique keys exercise full preallocation. + let rows = 300_000; + let on = vec![Arc::new(Column::new("key", 0)) as PhysicalExprRef]; + for distinct_keys in [20_000, rows] { + let batch = sampled_string_batch(rows, |row| { + Some(format!("key_{}", row % distinct_keys)) + })?; + let mut hashes = vec![0; rows]; + create_hashes(batch.columns(), HASH_JOIN_SEED.random_state(), &mut hashes)?; + assert_eq!( + check::(&batch, &on, &hashes, distinct_keys)?, + check::(&batch, &on, &hashes, distinct_keys)?, + ); + } + Ok(()) +} + +#[test] +fn sampled_hash_build_outgrown_hint_preserves_retained_memory() -> Result<()> { + let rows = 1_000_000; + let hot_keys = 30_000; + let unique_rows = 77_000; + let final_duplicate_rows = HASH_BUILD_CHUNK_ROWS; + let distinct_keys = hot_keys + unique_rows; + let batch = sampled_string_batch(rows, |row| { + // The builder visits this batch backwards: repeated keys arrive first, + // then 77K unique keys, then one duplicate chunk. That last chunk trips + // the conservative growth guard even though the distinct count is fixed. + let key = if (final_duplicate_rows..final_duplicate_rows + unique_rows) + .contains(&row) + { + hot_keys + row - final_duplicate_rows + } else { + row % hot_keys + }; + Some(format!("key_{key}")) + })?; + let column = Arc::new(Column::new("key", 0)) as PhysicalExprRef; + let sampled_on = vec![Arc::clone(&column)]; + // Repeating the same join key preserves its equivalence classes and chunk + // size but disables sampling, exercising the original growth policy. + let original_on = vec![Arc::clone(&column), column]; + let fixed_bytes = size_of::(); + let full_table_estimate = estimate_memory_size::<(u64, u32)>(rows, fixed_bytes)?; + let compact_table_estimate = + estimate_memory_size::<(u64, u32)>(distinct_keys, fixed_bytes)?; + // The original full allocation and its final compact replacement fit. + // During sampled growth, keeping the larger partial table and hash scratch + // live makes the full allocation fail by almost one scratch buffer. + let limit = fixed_bytes + + rows * size_of::() + + full_table_estimate + + compact_table_estimate; + let pool: Arc = Arc::new(GreedyMemoryPool::new(limit)); + let reservation = MemoryConsumer::new("retained-memory sample hint").register(&pool); + let (hint, _) = sampled_capacity( + std::slice::from_ref(&batch), + &sampled_on, + rows, + HASH_JOIN_SEED.random_state(), + &reservation, + ) + .expect("large flat string keys can be sampled"); + let hinted_entries = hint.saturating_add(HASH_BUILD_CHUNK_ROWS).min(rows); + assert_eq!( + estimate_memory_size::<(u64, u32)>(hinted_entries, fixed_bytes)?, + compact_table_estimate, + "fixture must initially select the table later outgrown by 107K keys: hint={hint}, entries={hinted_entries}", + ); + assert!( + sampled_capacity( + std::slice::from_ref(&batch), + &original_on, + rows, + HASH_JOIN_SEED.random_state(), + &reservation, + ) + .is_none() + ); + drop(reservation); + assert_eq!(pool.reserved(), 0); + + let build = |on: &[PhysicalExprRef]| -> Result<(usize, usize, usize, usize)> { + let recording = Arc::new(PeakRecordingPool::new(Arc::new( + GreedyMemoryPool::new(limit), + ))); + let pool: Arc = Arc::clone(&recording) as _; + let reservation = MemoryConsumer::new("outgrown sample hint").register(&pool); + let mut peak = 0; + let (table, next) = build_compact_hash_map::( + std::slice::from_ref(&batch), + on, + rows, + HASH_JOIN_SEED.random_state(), + NullEquality::NullEqualsNothing, + &reservation, + &mut peak, + )?; + assert_eq!(peak, recording.peak_reserved()); + assert_eq!( + reservation.size(), + fixed_bytes + rows * size_of::() + table.allocation_size() + ); + let summary = (table.len(), table.capacity(), table.allocation_size(), peak); + drop((table, next, reservation)); + assert_eq!(pool.reserved(), 0); + Ok(summary) + }; + let original = build(&original_on)?; + let sampled = build(&sampled_on)?; + assert_eq!(original.0, sampled.0); + if sampled.0 > 1 { + assert_eq!(sampled.0, distinct_keys); + assert!( + sampled.2 <= original.2, + "a sampled hint must not leave more retained memory after denied full growth: original={original:?}, sampled={sampled:?}", + ); + } + Ok(()) +} diff --git a/datafusion/physical-plan/src/joins/hash_join/exec.rs b/datafusion/physical-plan/src/joins/hash_join/exec.rs index 7b9e701119ef4..e3b89ffd964e6 100644 --- a/datafusion/physical-plan/src/joins/hash_join/exec.rs +++ b/datafusion/physical-plan/src/joins/hash_join/exec.rs @@ -31,6 +31,7 @@ use crate::filter_pushdown::{ }; use crate::joins::Map; use crate::joins::array_map::ArrayMap; +use crate::joins::hash_join::compact_hash_map::build_compact_hash_map; use crate::joins::hash_join::inlist_builder::build_struct_inlist_values; use crate::joins::hash_join::probe_completion::{ProbeCompletion, ProbeSideSummary}; use crate::joins::hash_join::shared_bounds::{ @@ -2921,6 +2922,8 @@ async fn collect_left_input( _ => None, }; + let mut hash_peak = 0; + let memory_before_hash = metrics.build_mem_used.value(); let (join_hash_map, batch, left_values) = if let Some((array_map, batch, left_value)) = try_create_array_map( bounds.as_ref(), @@ -2937,38 +2940,34 @@ async fn collect_left_input( (Map::ArrayMap(array_map), batch, left_value) } else { - // Estimation of memory size, required for hashtable, prior to allocation. - // Final result can be verified using `RawTable.allocation_info()` - // Use `u32` indices for the JoinHashMap when num_rows ≤ u32::MAX, otherwise use the - // `u64` indice variant - // Arc is used instead of Box to allow sharing with SharedBuildAccumulator for hash map pushdown - let mut hashmap = new_join_hashmap(num_rows, &mut reservation, &metrics)?; - - let mut hashes_buffer = Vec::new(); - let mut offset = 0; - - let batches_iter = batches.iter().rev(); - - // Updating hashmap starting from the last batch - for batch in batches_iter.clone() { - hashes_buffer.clear(); - hashes_buffer.resize(batch.num_rows(), 0); - update_hash( + let before = reservation.size(); + let hashmap: Box = if num_rows > u32::MAX as usize { + let (table, next) = build_compact_hash_map::( + &batches, &on_left, - batch, - &mut *hashmap, - offset, + num_rows, &random_state, - &mut hashes_buffer, - 0, - true, null_equality, + &reservation, + &mut hash_peak, )?; - offset += batch.num_rows(); - } + Box::new(JoinHashMapU64::new(table, next)) + } else { + let (table, next) = build_compact_hash_map::( + &batches, + &on_left, + num_rows, + &random_state, + null_equality, + &reservation, + &mut hash_peak, + )?; + Box::new(JoinHashMapU32::new(table, next)) + }; + metrics.build_mem_used.add(reservation.size() - before); // Merge all batches into a single batch, so we can directly index into the arrays - let batch = concat_batches(&schema, batches_iter.clone())?; + let batch = concat_batches(&schema, batches.iter().rev())?; let left_values = evaluate_expressions_to_arrays(&on_left, &batch)?; @@ -3104,6 +3103,12 @@ async fn collect_left_input( && !left_values.is_empty() && left_values[0].logical_null_count() > 0; + // Hash growth and scratch preceded the bitmap and auxiliary allocations. + // Report their peak without adding allocations that were not live together. + metrics + .build_mem_used + .set_max(memory_before_hash + hash_peak); + let data = JoinLeftData { map, null_aware_mark_scope_map, diff --git a/datafusion/physical-plan/src/joins/hash_join/mod.rs b/datafusion/physical-plan/src/joins/hash_join/mod.rs index 7c1b0d76c2f73..0ec33f4a4f7b9 100644 --- a/datafusion/physical-plan/src/joins/hash_join/mod.rs +++ b/datafusion/physical-plan/src/joins/hash_join/mod.rs @@ -20,6 +20,7 @@ pub use exec::{HashJoinExec, HashJoinExecBuilder}; pub use partitioned_hash_eval::{HashExpr, HashTableLookupExpr, SeededRandomState}; +mod compact_hash_map; mod exec; mod inlist_builder; mod partitioned_hash_eval; diff --git a/datafusion/physical-plan/src/joins/join_hash_map.rs b/datafusion/physical-plan/src/joins/join_hash_map.rs index 454cc916aeb12..4becb348d3396 100644 --- a/datafusion/physical-plan/src/joins/join_hash_map.rs +++ b/datafusion/physical-plan/src/joins/join_hash_map.rs @@ -30,8 +30,8 @@ use hashbrown::hash_table::Entry::{Occupied, Vacant}; /// Maps a `u64` hash value based on the build side ["on" values] to a list of indices with this key's value. /// -/// By allocating a `HashMap` with capacity for *at least* the number of rows for entries at the build side, -/// we make sure that we don't have to re-hash the hashmap, which needs access to the key (the hash in this case) value. +/// The lookup table needs one entry per distinct hash. The separate row-index +/// chain retains every build row, including rows with duplicate keys. /// /// E.g. 1 -> [3, 6, 8] indicates that the column values map to rows 3, 6 and 8 for hash value 1 /// As the key is a hash value, we need to check possible hash collisions in the probe stage @@ -149,7 +149,6 @@ pub struct JoinHashMapU32 { } impl JoinHashMapU32 { - #[cfg(test)] pub(crate) fn new(map: HashTable<(u64, u32)>, next: Vec) -> Self { Self { map, next } } @@ -229,7 +228,6 @@ pub struct JoinHashMapU64 { } impl JoinHashMapU64 { - #[cfg(test)] pub(crate) fn new(map: HashTable<(u64, u64)>, next: Vec) -> Self { Self { map, next } } @@ -312,6 +310,18 @@ pub fn update_from_iter<'a, T>( ) where T: Copy + TryFrom + PartialOrd, >::Error: Debug, +{ + update_from_iter_inner(map, next, iter, deleted_offset); +} + +pub(crate) fn update_from_iter_inner<'a, T>( + map: &mut HashTable<(u64, T)>, + next: &mut [T], + iter: impl Iterator + Send + 'a, + deleted_offset: usize, +) where + T: Copy + TryFrom + PartialOrd, + >::Error: Debug, { for (row, &hash_value) in iter { let entry = map.entry(