Skip to content

Commit 1bdf8bd

Browse files
adriangbclaude
andcommitted
fix(hash-join): measure only the work of a probe row without a match
A row that a dynamic filter removes is a probe row without a match. Its saving is the work that such a row gets: the evaluation and the hashes of the join keys and the hash table lookup. The check of the candidates (`equal_rows_arr`) and the output indices are work for the matches only. While the filter is on, most probe rows that reach the join are matches, thus this work made the measured saving too large. PR: #25681 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
1 parent 0b05ba2 commit 1bdf8bd

2 files changed

Lines changed: 69 additions & 25 deletions

File tree

‎datafusion/physical-expr/src/filter_stats.rs‎

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -170,9 +170,13 @@ impl FilterCost {
170170
/// probe row and looks them up in its hash table. A row that the dynamic
171171
/// filter of the join removes in the scan does not get this work, thus this
172172
/// is the saving of a removed row. The work that depends on a match (the
173-
/// output of a matched row) is not in it: the filter does not remove
174-
/// matched rows. The producer of a dynamic filter measures it and the
175-
/// consumers of the filter read it (see
173+
/// check of the candidates of the lookup and the output of a matched row)
174+
/// is not in it: the filter removes only rows without a match. While the
175+
/// filter is on, most rows that reach the producer are matches, thus work
176+
/// that only matches get would make the saving too large (TPC-DS SF1 Q31:
177+
/// 6 ns for each probe row with the check, 0.1 to 1.2 ns without it). The
178+
/// producer of a dynamic filter measures it and the consumers of the filter
179+
/// read it (see
176180
/// [`DynamicFilterPhysicalExpr::removed_row_work`]).
177181
///
178182
/// Lock-free.

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

Lines changed: 62 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -389,8 +389,10 @@ pub(super) struct HashJoinStream {
389389
/// batch and the time of the work that it does for a probe row whether
390390
/// or not the row matches: the evaluation and the hashes of the join
391391
/// keys and the hash table lookup. A row that the dynamic filter removes
392-
/// before the join does not get this work. Empty without dynamic filter
393-
/// pushdown.
392+
/// before the join does not get this work. The check of the candidates
393+
/// of the lookup and the output are not in it: only matches get them,
394+
/// and while the filter is on, most probe rows are matches. Empty
395+
/// without dynamic filter pushdown.
394396
removed_row_work: Vec<Arc<RemovedRowWork>>,
395397
/// Scratch space for probe indices during hash lookup
396398
probe_indices_buffer: Vec<u32>,
@@ -425,6 +427,18 @@ impl RecordBatchStream for HashJoinStream {
425427
}
426428
}
427429

430+
/// Records the time since `start` in `works` (see the `removed_row_work`
431+
/// field of [`HashJoinStream`]), without rows: the rows of a probe batch are
432+
/// recorded once, when the batch arrives.
433+
fn record_lookup_work(works: &[Arc<RemovedRowWork>], start: Option<Instant>) {
434+
if let Some(start) = start {
435+
let nanos = duration_nanos(start.elapsed());
436+
for work in works {
437+
work.record(0, nanos);
438+
}
439+
}
440+
}
441+
428442
/// Executes lookups by hash against JoinHashMap and resolves potential
429443
/// hash collisions.
430444
/// Returns build/probe indices satisfying the equality condition, along with
@@ -494,7 +508,27 @@ pub(super) fn lookup_join_hashmap(
494508
probe_indices_buffer,
495509
build_indices_buffer,
496510
);
511+
let (build_indices, probe_indices) = equal_candidates(
512+
build_side_values,
513+
probe_side_values,
514+
null_equality,
515+
probe_indices_buffer,
516+
build_indices_buffer,
517+
)?;
518+
Ok((build_indices, probe_indices, next_offset))
519+
}
497520

521+
/// Keeps the candidates of a hash table lookup (the indices in
522+
/// `probe_indices_buffer` and `build_indices_buffer`) whose join keys are
523+
/// equal, and returns their indices. Only the probe rows with a candidate
524+
/// (a match or a hash collision) get this work.
525+
fn equal_candidates(
526+
build_side_values: &[ArrayRef],
527+
probe_side_values: &[ArrayRef],
528+
null_equality: NullEquality,
529+
probe_indices_buffer: &mut Vec<u32>,
530+
build_indices_buffer: &mut Vec<u64>,
531+
) -> Result<(UInt64Array, UInt32Array)> {
498532
let build_indices_unfiltered: UInt64Array =
499533
std::mem::take(build_indices_buffer).into();
500534
let probe_indices_unfiltered: UInt32Array =
@@ -514,7 +548,7 @@ pub(super) fn lookup_join_hashmap(
514548
*build_indices_buffer = build_indices_unfiltered.into_parts().1.into();
515549
*probe_indices_buffer = probe_indices_unfiltered.into_parts().1.into();
516550

517-
Ok((build_indices, probe_indices, next_offset))
551+
Ok((build_indices, probe_indices))
518552
}
519553

520554
/// Counts the number of distinct elements in the input array.
@@ -916,22 +950,33 @@ impl HashJoinStream {
916950
return Ok(StatefulStreamResult::Continue);
917951
}
918952

919-
// get the matched by join keys indices
953+
// get the matched by join keys indices. The lookup of the hashes is
954+
// the work that every probe row gets, and a dynamic filter saves it
955+
// for each row that it removes: its time goes to `removed_row_work`.
956+
// The check and the output of the candidates is work for the matches
957+
// only, which a dynamic filter does not remove.
920958
let work_start = (!self.removed_row_work.is_empty()).then(Instant::now);
921959
let (left_indices, right_indices, next_offset) = match build_side.left_data.map()
922960
{
923-
Map::HashMap(map) => lookup_join_hashmap(
924-
map.as_ref(),
925-
build_side.left_data.values(),
926-
&state.values,
927-
self.null_equality,
928-
&self.hashes_buffer,
929-
state.valid_keys.as_ref(),
930-
self.batch_size,
931-
state.offset,
932-
&mut self.probe_indices_buffer,
933-
&mut self.build_indices_buffer,
934-
)?,
961+
Map::HashMap(map) => {
962+
let next_offset = map.get_matched_indices_with_limit_offset(
963+
&self.hashes_buffer,
964+
state.valid_keys.as_ref(),
965+
self.batch_size,
966+
state.offset,
967+
&mut self.probe_indices_buffer,
968+
&mut self.build_indices_buffer,
969+
);
970+
record_lookup_work(&self.removed_row_work, work_start);
971+
let (left_indices, right_indices) = equal_candidates(
972+
build_side.left_data.values(),
973+
&state.values,
974+
self.null_equality,
975+
&mut self.probe_indices_buffer,
976+
&mut self.build_indices_buffer,
977+
)?;
978+
(left_indices, right_indices, next_offset)
979+
}
935980
Map::ArrayMap(array_map) => {
936981
let next_offset = array_map.get_matched_indices_with_limit_offset(
937982
&state.values,
@@ -940,19 +985,14 @@ impl HashJoinStream {
940985
&mut self.probe_indices_buffer,
941986
&mut self.build_indices_buffer,
942987
)?;
988+
record_lookup_work(&self.removed_row_work, work_start);
943989
(
944990
UInt64Array::from(self.build_indices_buffer.clone()),
945991
UInt32Array::from(self.probe_indices_buffer.clone()),
946992
next_offset,
947993
)
948994
}
949995
};
950-
951-
if let Some(work_start) = work_start {
952-
for work in &self.removed_row_work {
953-
work.record(0, duration_nanos(work_start.elapsed()));
954-
}
955-
}
956996
let matched_probe_rows = state.count_new_matched_probe_rows(&right_indices);
957997

958998
self.join_metrics

0 commit comments

Comments
 (0)