Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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
522 changes: 282 additions & 240 deletions datafusion/core/tests/physical_optimizer/filter_pushdown.rs

Large diffs are not rendered by default.

117 changes: 96 additions & 21 deletions datafusion/physical-plan/src/joins/hash_join/exec.rs
Original file line number Diff line number Diff line change
Expand Up @@ -876,24 +876,31 @@ impl HashJoinExec {
return false;
}

// `preserve_file_partitions` can report Hash partitioning for Hive-style
// file groups, but those partitions are not actually hash-distributed.
// Partitioned dynamic filters rely on hash routing, so disable them in
// this mode to avoid incorrect results. Follow-up work: enable dynamic
// filtering for preserve_file_partitioned scans (issue #20195).
// `preserve_file_partitions` can report Hive-style file groups as Hash
// partitioned even though their partition indexes do not follow the
// hash router used by partitioned dynamic filters. Reject Hash inputs
// because the metadata cannot distinguish those scans from a real hash
// repartition. Compatible Range inputs remain safe because matching
// ordering and split points align each build filter with its probe
// partition. Other unsupported layouts are rejected.
// Follow-up work: enable dynamic filtering for preserve_file_partitioned scans (issue #20195).
// https://github.com/apache/datafusion/issues/20195
if config.optimizer.preserve_file_partitions > 0
&& self.mode == PartitionMode::Partitioned
&& matches!(
(
self.left.output_partitioning(),
self.right.output_partitioning()
),
(Partitioning::Hash(_, _), Partitioning::Hash(_, _))
)
{
return false;
}

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 @@ -909,6 +916,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 @@ -7089,25 +7104,26 @@ mod tests {
Ok(())
}

#[test]
fn test_partitioned_dynamic_filter_pushdown_rejects_range_partitioning() -> Result<()>
{
fn range_partitioned_dynamic_filter_test_join(
left_split: i32,
right_split: i32,
) -> Result<(HashJoinExec, JoinOn)> {
let (left_schema, right_schema, on) = build_schema_and_on()?;
let left_partitioning = Partitioning::Range(RangePartitioning::try_new(
[PhysicalSortExpr {
expr: Arc::clone(&on[0].0),
options: Default::default(),
}]
.into(),
vec![SplitPoint::new(vec![ScalarValue::Int32(Some(10))])],
vec![SplitPoint::new(vec![ScalarValue::Int32(Some(left_split))])],
)?);
let 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(10))])],
vec![SplitPoint::new(vec![ScalarValue::Int32(Some(right_split))])],
)?);
let left = Arc::new(PartitionedTestExec::try_new(
left_schema,
Expand All @@ -7118,25 +7134,84 @@ mod tests {
right_partitioning,
)?);

let mut session_config = SessionConfig::default();
session_config
.options_mut()
.optimizer
.enable_join_dynamic_filter_pushdown = true;

let join = HashJoinExec::try_new(
left,
right,
on,
on.clone(),
None,
&JoinType::Inner,
None,
PartitionMode::Partitioned,
NullEquality::NullEqualsNothing,
false,
)?;
Ok((join, on))
}

assert!(!join.allow_join_dynamic_filter_pushdown(session_config.options()));
fn with_hash_partitioned_children(
join: &HashJoinExec,
on: &JoinOn,
) -> Result<HashJoinExec> {
join.builder()
.with_new_children(vec![
Arc::new(PartitionedTestExec::try_new(
join.left().schema(),
Partitioning::Hash(vec![Arc::clone(&on[0].0)], 2),
)?),
Arc::new(PartitionedTestExec::try_new(
join.right().schema(),
Partitioning::Hash(vec![Arc::clone(&on[0].1)], 2),
)?),
])?
.build()
}

#[test]
fn test_partitioned_dynamic_filter_pushdown_allows_supported_partitioning()
-> Result<()> {
let (range_join, on) = range_partitioned_dynamic_filter_test_join(10, 10)?;
let hash_join = with_hash_partitioned_children(&range_join, &on)?;
let mut session_config = SessionConfig::default();
session_config
.options_mut()
.optimizer
.enable_join_dynamic_filter_pushdown = true;

assert!(range_join.allow_join_dynamic_filter_pushdown(session_config.options()));
assert!(hash_join.allow_join_dynamic_filter_pushdown(session_config.options()));

session_config
.options_mut()
.optimizer
.preserve_file_partitions = 1;
assert!(range_join.allow_join_dynamic_filter_pushdown(session_config.options()));

Ok(())
}

#[test]
fn test_partitioned_dynamic_filter_pushdown_rejects_unsupported_partitioning()
-> Result<()> {
let (range_join, on) = range_partitioned_dynamic_filter_test_join(10, 10)?;
let hash_join = with_hash_partitioned_children(&range_join, &on)?;
let (mismatched_range_join, _) =
range_partitioned_dynamic_filter_test_join(10, 11)?;
let mut session_config = SessionConfig::default();
session_config
.options_mut()
.optimizer
.enable_join_dynamic_filter_pushdown = true;

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

session_config
.options_mut()
.optimizer
.preserve_file_partitions = 1;
assert!(!hash_join.allow_join_dynamic_filter_pushdown(session_config.options()));

Ok(())
}
Expand Down
Loading
Loading