Skip to content

Commit 16446e3

Browse files
author
B Vadlamani
committed
optimizations_roaring_memebership_check
1 parent 8c2e5c1 commit 16446e3

4 files changed

Lines changed: 47 additions & 11 deletions

File tree

datafusion/physical-plan/src/joins/hash_join/exec.rs

Lines changed: 10 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,8 @@ use crate::filter_pushdown::{
3131
ChildFilterDescription, ChildPushdownResult, FilterDescription, FilterPushdownPhase,
3232
FilterPushdownPropagation,
3333
};
34-
use crate::joins::Map;
34+
use crate::joins::{Map, RoaringMapData};
35+
use roaring::RoaringBitmap;
3536
use crate::joins::array_map::ArrayMap;
3637
use crate::joins::hash_join::inlist_builder::build_struct_inlist_values;
3738
use crate::joins::hash_join::shared_bounds::{
@@ -1994,7 +1995,7 @@ async fn collect_left_input(
19941995
let batch = concat_batches(&schema, &batches)?;
19951996
let left_values = evaluate_expressions_to_arrays(&on_left, &batch)?;
19961997
let key_col = &left_values[0];
1997-
let bitmap = match key_col.data_type() {
1998+
let bitmap: RoaringBitmap = match key_col.data_type() {
19981999
DataType::UInt32 => key_col
19992000
.as_any()
20002001
.downcast_ref::<UInt32Array>()
@@ -2013,7 +2014,13 @@ async fn collect_left_input(
20132014
_ => return internal_err!("unsupported data type to build bitmap"),
20142015
};
20152016
roaring_map_created_count.add(1);
2016-
(Map::RoaringMap(bitmap), batch, left_values)
2017+
let min = bitmap.min().unwrap_or(u32::MAX);
2018+
let max = bitmap.max().unwrap_or(u32::MIN);
2019+
(
2020+
Map::RoaringMap(RoaringMapData { bitmap, min, max }),
2021+
batch,
2022+
left_values,
2023+
)
20172024
} else if let Some((array_map, batch, left_value)) = try_create_array_map(
20182025
&bounds,
20192026
&schema,

datafusion/physical-plan/src/joins/hash_join/partitioned_hash_eval.rs

Lines changed: 11 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -336,7 +336,7 @@ impl PhysicalExpr for HashTableLookupExpr {
336336
let array = map.contain_keys(&join_keys)?;
337337
Ok(ColumnarValue::Array(Arc::new(array)))
338338
}
339-
Map::RoaringMap(bitmap) => {
339+
Map::RoaringMap(data) => {
340340
// For roaring bitmap, check membership directly
341341
// Roaring only supports u32 values, so we handle Int32 and UInt32 types
342342
if join_keys.len() != 1 {
@@ -346,6 +346,12 @@ impl PhysicalExpr for HashTableLookupExpr {
346346
);
347347
}
348348
let key_col = &join_keys[0];
349+
// Bounds were cached at build time so out-of-range probes
350+
// short-circuit before the binary-search inside
351+
// `RoaringBitmap::contains`.
352+
let min = data.min;
353+
let max = data.max;
354+
let bitmap = &data.bitmap;
349355
let contains: BooleanArray = match key_col.data_type() {
350356
DataType::Int32 => {
351357
let arr = key_col
@@ -356,7 +362,8 @@ impl PhysicalExpr for HashTableLookupExpr {
356362
if arr.is_null(i) {
357363
return false;
358364
}
359-
bitmap.contains(arr.value(i) as u32)
365+
let v = arr.value(i) as u32;
366+
v >= min && v <= max && bitmap.contains(v)
360367
});
361368
BooleanArray::new(buffer, None)
362369
}
@@ -369,7 +376,8 @@ impl PhysicalExpr for HashTableLookupExpr {
369376
if arr.is_null(i) {
370377
return false;
371378
}
372-
bitmap.contains(arr.value(i))
379+
let v = arr.value(i);
380+
v >= min && v <= max && bitmap.contains(v)
373381
});
374382
BooleanArray::new(buffer, None)
375383
}

datafusion/physical-plan/src/joins/hash_join/stream.rs

Lines changed: 15 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -732,17 +732,26 @@ impl HashJoinStream {
732732
next_offset,
733733
)
734734
}
735-
Map::RoaringMap(bitmap) => {
735+
Map::RoaringMap(data) => {
736736
let key_col = &state.values[0];
737737
let is_semi = matches!(self.join_type, JoinType::RightSemi);
738+
// Bounds were cached at build time so out-of-range probes
739+
// short-circuit before the binary-search inside
740+
// `RoaringBitmap::contains`.
741+
let min = data.min;
742+
let max = data.max;
743+
let bitmap = &data.bitmap;
738744
let mask: BooleanArray = match key_col.data_type() {
739745
DataType::Int32 => {
740746
let arr = key_col.as_any().downcast_ref::<Int32Array>().unwrap();
741747
let buffer = BooleanBuffer::collect_bool(arr.len(), |i| {
742748
if arr.is_null(i) {
743749
return !is_semi;
744750
}
745-
bitmap.contains(arr.value(i) as u32) == is_semi
751+
let v = arr.value(i) as u32;
752+
let in_set =
753+
v >= min && v <= max && bitmap.contains(v);
754+
in_set == is_semi
746755
});
747756
BooleanArray::new(buffer, None)
748757
}
@@ -752,7 +761,10 @@ impl HashJoinStream {
752761
if arr.is_null(i) {
753762
return !is_semi;
754763
}
755-
bitmap.contains(arr.value(i)) == is_semi
764+
let v = arr.value(i);
765+
let in_set =
766+
v >= min && v <= max && bitmap.contains(v);
767+
in_set == is_semi
756768
});
757769
BooleanArray::new(buffer, None)
758770
}

datafusion/physical-plan/src/joins/mod.rs

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -51,11 +51,20 @@ pub mod join_hash_map;
5151
use array_map::ArrayMap;
5252
use utils::JoinHashMapType;
5353

54+
pub struct RoaringMapData {
55+
pub bitmap: RoaringBitmap,
56+
/// Smallest key in `bitmap`. Cached at build time so probe loops can
57+
/// short-circuit out-of-range keys without invoking `RoaringBitmap::contains`.
58+
pub min: u32,
59+
/// Largest key in `bitmap`. See [`Self::min`].
60+
pub max: u32,
61+
}
62+
5463
pub enum Map {
5564
HashMap(Box<dyn JoinHashMapType>),
5665
ArrayMap(ArrayMap),
5766
// optimized path for single int join keys
58-
RoaringMap(RoaringBitmap),
67+
RoaringMap(RoaringMapData),
5968
}
6069

6170
impl Map {
@@ -64,7 +73,7 @@ impl Map {
6473
match self {
6574
Map::HashMap(map) => map.len(),
6675
Map::ArrayMap(array_map) => array_map.num_of_distinct_key(),
67-
Map::RoaringMap(bitmap) => bitmap.len() as usize,
76+
Map::RoaringMap(data) => data.bitmap.len() as usize,
6877
}
6978
}
7079

0 commit comments

Comments
 (0)