Skip to content

Commit 86c0f30

Browse files
author
B Vadlamani
committed
null_handling
1 parent cc6607d commit 86c0f30

2 files changed

Lines changed: 18 additions & 34 deletions

File tree

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

Lines changed: 4 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -1993,9 +1993,7 @@ async fn collect_left_input(
19931993
};
19941994

19951995
let (join_hash_map, batch, left_values) =
1996-
if config.execution.perfect_hash_join_small_build_threshold > 0
1997-
&& should_use_roaring_bitmap(&join_type, &on_left, &schema, has_filter)
1998-
{
1996+
if should_use_roaring_bitmap(&join_type, &on_left, &schema, has_filter) {
19991997
let batch = concat_batches(&schema, &batches)?;
20001998
let left_values = evaluate_expressions_to_arrays(&on_left, &batch)?;
20011999
let key_col = &left_values[0];
@@ -2004,17 +2002,16 @@ async fn collect_left_input(
20042002
.as_any()
20052003
.downcast_ref::<UInt32Array>()
20062004
.unwrap()
2007-
.values()
20082005
.iter()
2009-
.map(|&v| v)
2006+
.flatten()
20102007
.collect(),
20112008
DataType::Int32 => key_col
20122009
.as_any()
20132010
.downcast_ref::<Int32Array>()
20142011
.unwrap()
2015-
.values()
20162012
.iter()
2017-
.map(|&v| v as u32)
2013+
.flatten()
2014+
.map(|v| v as u32)
20182015
.collect(),
20192016
_ => return internal_err!("unsupported data type to build bitmap"),
20202017
};

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

Lines changed: 14 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -46,7 +46,8 @@ use crate::{
4646
},
4747
};
4848

49-
use arrow::array::{Array, ArrayRef, Int32Array, UInt32Array, UInt64Array};
49+
use arrow::array::{Array, ArrayRef, BooleanArray, Int32Array, UInt32Array, UInt64Array};
50+
use arrow::compute::filter_record_batch;
5051
use arrow::datatypes::{Schema, SchemaRef};
5152
use arrow::record_batch::RecordBatch;
5253
use arrow_schema::DataType;
@@ -727,44 +728,30 @@ impl HashJoinStream {
727728
Map::RoaringMap(bitmap) => {
728729
let key_col = &state.values[0];
729730
let is_semi = matches!(self.join_type, JoinType::RightSemi);
730-
let right_indices = match key_col.data_type() {
731+
let mask: BooleanArray = match key_col.data_type() {
731732
DataType::Int32 => {
732733
let arr = key_col.as_any().downcast_ref::<Int32Array>().unwrap();
733-
arr.values()
734-
.iter()
735-
.enumerate()
736-
.filter_map(|(i, v)| {
737-
let contains = bitmap.contains(*v as u32);
738-
let emit = if is_semi { contains } else { !contains };
739-
emit.then_some(i as u32)
734+
arr.iter()
735+
.map(|v| match v {
736+
Some(v) => bitmap.contains(v as u32) == is_semi,
737+
None => !is_semi,
740738
})
741-
.collect::<Vec<u32>>()
739+
.collect()
742740
}
743741
DataType::UInt32 => {
744742
let arr = key_col.as_any().downcast_ref::<UInt32Array>().unwrap();
745-
arr.values()
746-
.iter()
747-
.enumerate()
748-
.filter_map(|(i, v)| {
749-
let contains = bitmap.contains(*v);
750-
let emit = if is_semi { contains } else { !contains };
751-
emit.then_some(i as u32)
743+
arr.iter()
744+
.map(|v| match v {
745+
Some(v) => bitmap.contains(v) == is_semi,
746+
None => !is_semi,
752747
})
753-
.collect::<Vec<u32>>()
748+
.collect()
754749
}
755750
_ => {
756751
return internal_err!("unsupported data type for roaring bitmap");
757752
}
758753
};
759-
let indices = UInt32Array::from(right_indices);
760-
let columns: Vec<ArrayRef> = state
761-
.batch
762-
.columns()
763-
.iter()
764-
.map(|col| arrow::compute::take(col, &indices, None))
765-
.collect::<Result<Vec<_>, _>>()?;
766-
767-
let batch = RecordBatch::try_new(self.schema.clone(), columns)?;
754+
let batch = filter_record_batch(&state.batch, &mask)?;
768755
self.output_buffer.push_batch(batch)?;
769756
self.state = HashJoinStreamState::FetchProbeBatch;
770757
return Ok(StatefulStreamResult::Continue);

0 commit comments

Comments
 (0)