Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
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
69 changes: 66 additions & 3 deletions datafusion/core/src/physical_planner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,8 @@ use crate::physical_plan::explain::ExplainExec;
use crate::physical_plan::filter::FilterExecBuilder;
use crate::physical_plan::joins::utils as join_utils;
use crate::physical_plan::joins::{
CrossJoinExec, HashJoinExec, NestedLoopJoinExec, PartitionMode, SortMergeJoinExec,
AsOfJoinExec, AsOfMatchExpr, CrossJoinExec, HashJoinExec, NestedLoopJoinExec,
PartitionMode, SortMergeJoinExec,
};
use crate::physical_plan::limit::{GlobalLimitExec, LocalLimitExec};
use crate::physical_plan::projection::{ProjectionExec, ProjectionExpr};
Expand Down Expand Up @@ -93,8 +94,8 @@ use datafusion_expr::physical_planning_context::{
use datafusion_expr::utils::{expr_to_columns, split_conjunction};
use datafusion_expr::{
Analyze, BinaryExpr, DescribeTable, DmlStatement, Explain, ExplainFormat, Extension,
FetchType, Filter, JoinType, Operator, RecursiveQuery, SkipType, StringifiedPlan,
WindowFrame, WindowFrameBound, WriteOp,
FetchType, Filter, JoinConstraint, JoinType, Operator, RecursiveQuery, SkipType,
StringifiedPlan, WindowFrame, WindowFrameBound, WriteOp,
};
use datafusion_physical_expr::aggregate::{
AggregateFunctionExpr, LoweredAggregate, LoweredAggregateBuilder,
Expand Down Expand Up @@ -1860,6 +1861,67 @@ impl DefaultPhysicalPlanner {
join
}
}
LogicalPlan::AsOfJoin(join) => {
let [physical_left, physical_right] = children.two()?;
let join_on = join
.on
.iter()
.map(|(left, right)| {
Ok((
create_physical_expr(
left,
join.left.schema(),
execution_props,
planning_ctx,
)?,
create_physical_expr(
right,
join.right.schema(),
execution_props,
planning_ctx,
)?,
))
})
.collect::<Result<join_utils::JoinOn>>()?;
let match_condition = AsOfMatchExpr::new(
create_physical_expr(
&join.match_condition.left,
join.left.schema(),
execution_props,
planning_ctx,
)?,
join.match_condition.op,
create_physical_expr(
&join.match_condition.right,
join.right.schema(),
execution_props,
planning_ctx,
)?,
);
let omitted_right = if join.join_constraint == JoinConstraint::Using {
join.on
.iter()
.map(|(_, right)| {
let column = right.get_as_join_column().ok_or_else(|| {
internal_datafusion_err!("ASOF USING key is not a column")
})?;
join.right.schema().index_of_column(column)
})
.collect::<Result<HashSet<_>>>()?
} else {
HashSet::new()
};
let right_output_indices = (0..join.right.schema().fields().len())
.filter(|index| !omitted_right.contains(index))
.collect();
Arc::new(AsOfJoinExec::try_new(
physical_left,
physical_right,
join_on,
match_condition,
right_output_indices,
)?)
}
LogicalPlan::RecursiveQuery(RecursiveQuery {
name,
is_distinct,
Expand Down Expand Up @@ -2359,6 +2421,7 @@ fn extract_dml_filters(
| LogicalPlan::Sort(_)
| LogicalPlan::Union(_)
| LogicalPlan::Join(_)
| LogicalPlan::AsOfJoin(_)
| LogicalPlan::Repartition(_)
| LogicalPlan::Aggregate(_)
| LogicalPlan::Window(_)
Expand Down
42 changes: 42 additions & 0 deletions datafusion/core/tests/physical_optimizer/filter_pushdown.rs
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,7 @@ use datafusion_physical_plan::{
coalesce_partitions::CoalescePartitionsExec,
collect,
filter::{FilterExec, FilterExecBuilder},
joins::{AsOfJoinExec, AsOfMatchExpr},
projection::ProjectionExec,
repartition::RepartitionExec,
sorts::sort::SortExec,
Expand Down Expand Up @@ -123,6 +124,47 @@ fn test_pushdown_volatile_functions_not_allowed() {
);
}

#[test]
fn test_asof_join_pushes_only_left_filters() {
let left = TestScanBuilder::new(schema()).with_support(true).build();
let right = TestScanBuilder::new(schema()).with_support(true).build();
let join = Arc::new(
AsOfJoinExec::try_new(
left,
right,
vec![(col("a", &schema()).unwrap(), col("a", &schema()).unwrap())],
AsOfMatchExpr::new(
col("b", &schema()).unwrap(),
Operator::GtEq,
col("b", &schema()).unwrap(),
),
vec![0, 1, 2],
)
.unwrap(),
);
let left_filter = Arc::new(BinaryExpr::new(
Arc::new(Column::new("c", 2)),
Operator::Gt,
Arc::new(Literal::new(ScalarValue::Float64(Some(0.0)))),
)) as Arc<dyn PhysicalExpr>;
let right_filter = Arc::new(BinaryExpr::new(
Arc::new(Column::new("c", 5)),
Operator::Gt,
Arc::new(Literal::new(ScalarValue::Float64(Some(0.0)))),
)) as Arc<dyn PhysicalExpr>;
let predicate = Arc::new(BinaryExpr::new(left_filter, Operator::And, right_filter));
let plan = Arc::new(FilterExec::try_new(predicate, join).unwrap());
let mut config = ConfigOptions::default();
config.execution.parquet.pushdown_filters = true;
let optimized = FilterPushdown::new().optimize(plan, &config).unwrap();
let formatted = format_plan_for_test(&optimized);

assert!(formatted.contains("FilterExec: c@5 > 0"), "{formatted}");
assert!(formatted.contains("AsOfJoinExec:"), "{formatted}");
assert!(formatted.contains("predicate=c@2 > 0"), "{formatted}");
assert_eq!(formatted.matches("predicate=").count(), 1, "{formatted}");
}

/// Show that we can use config options to determine how to do pushdown.
#[test]
fn test_pushdown_into_scan_with_config_options() {
Expand Down
Loading
Loading