Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 commits
Commits
Show all changes
25 commits
Select commit Hold shift + click to select a range
0c5a549
Enable filters for range-partitioned joins
peterxcli Jul 23, 2026
0b71b56
Merge branch 'main' of https://github.com/apache/datafusion into feat…
peterxcli Jul 23, 2026
a9167cb
revert header change
peterxcli Jul 23, 2026
4206b36
remove float zero test
peterxcli Jul 24, 2026
1a141f3
address review batch 1
peterxcli Jul 24, 2026
d6b843b
Allow dynamic filters for range-partitioned joins
peterxcli Jul 31, 2026
bf4e6a1
Merge remote-tracking branch 'upstream/main' into feat/hash-join-dyna…
peterxcli Jul 31, 2026
7d2589f
Merge branch 'main' into feat/hash-join-dynamic-filter-with-range-par…
peterxcli Jul 31, 2026
95ad3ca
jay's review
peterxcli Aug 4, 2026
dd7fefa
add diagram
peterxcli Aug 4, 2026
10ee4c7
jay's 2nd review
peterxcli Aug 5, 2026
18f9db3
Merge branch 'main' into feat/hash-join-dynamic-filter-with-range-par…
peterxcli Aug 5, 2026
3a9d253
dedup filter pushdown test
peterxcli Aug 8, 2026
6dae145
review for exec.rs
peterxcli Aug 8, 2026
7f84831
replace csv with pq and add slt
peterxcli Aug 8, 2026
fdb292f
Merge branch 'main' into feat/hash-join-dynamic-filter-with-range-par…
peterxcli Aug 9, 2026
4990a57
cargo fmt
peterxcli Aug 9, 2026
ee07f7b
remove duplicate imports
peterxcli Aug 10, 2026
8903354
CASE WHEN over range_split[...] -> range_partition PhysicalExpr
peterxcli Aug 10, 2026
37fb97b
Merge branch 'main' into feat/hash-join-dynamic-filter-with-range-par…
peterxcli Aug 10, 2026
eaf7f86
codegen
peterxcli Aug 10, 2026
01210bd
revert topk filter builder
peterxcli Aug 11, 2026
00ef9f0
assert_or_internal_err, top import and on_columns accessor
peterxcli Aug 11, 2026
f9f1f0d
check the range and sort properties
peterxcli Aug 11, 2026
7274de3
Merge branch 'main' into feat/hash-join-dynamic-filter-with-range-par…
peterxcli Aug 11, 2026
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
106 changes: 99 additions & 7 deletions datafusion/physical-plan/src/joins/hash_join/exec.rs
Original file line number Diff line number Diff line change
Expand Up @@ -886,9 +886,6 @@ impl HashJoinExec {
if self.mode == PartitionMode::Partitioned
&& !self.has_partitioned_dynamic_filter_routing()
{
// TODO: support partition-routed dynamic filters for compatible
// range co-partitioned joins.
// <https://github.com/apache/datafusion/issues/23376>.
return false;
}

Expand All @@ -904,6 +901,14 @@ impl HashJoinExec {
Partitioning::Hash(_, left_partition_count),
Partitioning::Hash(_, right_partition_count),
) => left_partition_count == right_partition_count,
(Partitioning::Range(_), Partitioning::Range(_)) => {
let children = [self.left.as_ref(), self.right.as_ref()];
matches!(
self.input_distribution_requirements()
.unsatisfied_co_partitioned_children(self.name(), &children),
Ok(unsatisfied) if unsatisfied.is_empty()
)
}
(left_partitioning, right_partitioning) => {
left_partitioning.partition_count() == 1
&& right_partitioning.partition_count() == 1
Expand Down Expand Up @@ -6759,8 +6764,7 @@ mod tests {
}

#[test]
fn test_partitioned_dynamic_filter_pushdown_rejects_range_partitioning() -> Result<()>
{
fn test_partitioned_dynamic_filter_pushdown_range_partitioning() -> Result<()> {
Comment thread
peterxcli marked this conversation as resolved.
Outdated
let (left_schema, right_schema, on) = build_schema_and_on()?;
let left_partitioning = Partitioning::Range(RangePartitioning::try_new(
[PhysicalSortExpr {
Expand All @@ -6783,7 +6787,7 @@ mod tests {
left_partitioning,
)?);
let right = Arc::new(PartitionedTestExec::try_new(
right_schema,
Arc::clone(&right_schema),
right_partitioning,
)?);

Expand All @@ -6794,8 +6798,35 @@ mod tests {
.enable_join_dynamic_filter_pushdown = true;

let join = HashJoinExec::try_new(
left,
Arc::clone(&left) as Arc<dyn ExecutionPlan>,
right,
on.clone(),
None,
&JoinType::Inner,
None,
PartitionMode::Partitioned,
NullEquality::NullEqualsNothing,
false,
)?;

assert!(join.allow_join_dynamic_filter_pushdown(session_config.options()));
Comment thread
peterxcli marked this conversation as resolved.
Outdated

let mismatched_right_partitioning =
Partitioning::Range(RangePartitioning::try_new(
[PhysicalSortExpr {
expr: Arc::clone(&on[0].1),
options: Default::default(),
}]
.into(),
vec![SplitPoint::new(vec![ScalarValue::Int32(Some(11))])],
)?);
let mismatched_right = Arc::new(PartitionedTestExec::try_new(
right_schema,
mismatched_right_partitioning,
)?);
let mismatched_join = HashJoinExec::try_new(
left,
mismatched_right,
on,
None,
&JoinType::Inner,
Expand All @@ -6805,6 +6836,67 @@ mod tests {
false,
)?;

assert!(
!mismatched_join.allow_join_dynamic_filter_pushdown(session_config.options())
);

Ok(())
}

#[test]
fn test_partitioned_dynamic_filter_pushdown_rejects_float_zero_split() -> Result<()> {
let left_schema = Arc::new(Schema::new(vec![Field::new(
"left_key",
DataType::Float64,
false,
)]));
let right_schema = Arc::new(Schema::new(vec![Field::new(
"right_key",
DataType::Float64,
false,
)]));
let left_key = Arc::new(Column::new("left_key", 0)) as PhysicalExprRef;
let right_key = Arc::new(Column::new("right_key", 0)) as PhysicalExprRef;
let split_points = vec![SplitPoint::new(vec![ScalarValue::Float64(Some(0.0))])];
let left = Arc::new(PartitionedTestExec::try_new(
left_schema,
Partitioning::Range(RangePartitioning::try_new(
[PhysicalSortExpr::new(
Arc::clone(&left_key),
Default::default(),
)]
.into(),
split_points.clone(),
)?),
)?);
let right = Arc::new(PartitionedTestExec::try_new(
right_schema,
Partitioning::Range(RangePartitioning::try_new(
[PhysicalSortExpr::new(
Arc::clone(&right_key),
Default::default(),
)]
.into(),
split_points,
)?),
)?);
let join = HashJoinExec::try_new(
left,
right,
vec![(left_key, right_key)],
None,
&JoinType::Inner,
None,
PartitionMode::Partitioned,
NullEquality::NullEqualsNothing,
false,
)?;
let mut session_config = SessionConfig::default();
session_config
.options_mut()
.optimizer
.enable_join_dynamic_filter_pushdown = true;

assert!(!join.allow_join_dynamic_filter_pushdown(session_config.options()));

Ok(())
Expand Down
Loading
Loading