Skip to content
Draft
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
86 changes: 78 additions & 8 deletions datafusion/common/src/hash_utils.rs
Original file line number Diff line number Diff line change
Expand Up @@ -321,15 +321,9 @@ fn hash_array_primitive<T>(

if array.null_count() == 0 {
if rehash {
for (hash, &value) in hashes_buffer.iter_mut().zip(array.values().iter()) {
let mut hasher = random_state.seeded_state(*hash).build_hasher();
value.hash_write(&mut hasher);
*hash = hasher.finish();
}
hash_prim_rehash_dense(array.values(), hashes_buffer, random_state);
} else {
for (hash, &value) in hashes_buffer.iter_mut().zip(array.values().iter()) {
*hash = value.hash_one(random_state);
}
hash_prim_fresh_dense(array.values(), hashes_buffer, random_state);
}
} else if rehash {
for i in array.nulls().unwrap().valid_indices() {
Expand All @@ -346,6 +340,82 @@ fn hash_array_primitive<T>(
}
}

/// Overwrite each slot with `value.hash_one(state)`, 16 rows per iteration.
/// Gathering + hashing a whole lane batch before storing keeps the independent
/// hashes in flight so the pipeline isn't serialized on store-to-load deps.
#[cfg(not(feature = "force_hash_collisions"))]
#[inline(always)]
fn hash_prim_fresh_dense<N, S>(values: &[N], hashes: &mut [u64], state: &S)
where
N: HashValue + Copy,
S: HashState,
{
debug_assert_eq!(values.len(), hashes.len());
const LANES: usize = 16;
let len = values.len();
let aligned_end = len & !(LANES - 1);
let mut row = 0;
while row < aligned_end {
unsafe {
let new_hashes: [u64; LANES] = std::array::from_fn(|lane| {
values.get_unchecked(row + lane).hash_one(state)
});
for lane in 0..LANES {
*hashes.get_unchecked_mut(row + lane) = new_hashes[lane];
}
}
row += LANES;
}
while row < len {
unsafe {
*hashes.get_unchecked_mut(row) = values.get_unchecked(row).hash_one(state);
}
row += 1;
}
}

/// Fold `prev_hash` + `value` into a new hash, 8 rows per iteration. 8 is
/// narrower than `hash_prim_fresh_dense` because each lane holds both `prev` and
/// `value` live; widening past 8 spilled registers in the multi-column case.
#[cfg(not(feature = "force_hash_collisions"))]
#[inline(always)]
fn hash_prim_rehash_dense<N, S>(values: &[N], hashes: &mut [u64], state: &S)
where
N: HashValue + Copy,
S: HashState,
{
debug_assert_eq!(values.len(), hashes.len());
const LANES: usize = 8;
let len = values.len();
let aligned_end = len & !(LANES - 1);
let mut row = 0;
while row < aligned_end {
unsafe {
let new_hashes: [u64; LANES] = std::array::from_fn(|lane| {
let value = *values.get_unchecked(row + lane);
let prev_hash = *hashes.get_unchecked(row + lane);
let mut hasher = state.seeded_state(prev_hash).build_hasher();
value.hash_write(&mut hasher);
hasher.finish()
});
for lane in 0..LANES {
*hashes.get_unchecked_mut(row + lane) = new_hashes[lane];
}
}
row += LANES;
}
while row < len {
unsafe {
let value = *values.get_unchecked(row);
let prev_hash = *hashes.get_unchecked(row);
let mut hasher = state.seeded_state(prev_hash).build_hasher();
value.hash_write(&mut hasher);
*hashes.get_unchecked_mut(row) = hasher.finish();
}
row += 1;
}
}

/// Hashes one array into the `hashes_buffer`
/// If `rehash==true` this combines the previous hash value in the buffer
/// with the new hash using `combine_hashes`
Expand Down
Loading