Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
e546de5
feat: add ASOF join physical operator
Xuanwo Jul 23, 2026
bd58e1d
feat: add ASOF join logical semantics
Xuanwo Jul 23, 2026
ed2ce92
feat: serialize ASOF join plans
Xuanwo Jul 23, 2026
68ac8ec
Merge remote-tracking branch 'origin/main' into xuanwo/asof-physical
Xuanwo Jul 23, 2026
c8c4813
Merge branch 'xuanwo/asof-physical' into xuanwo/asof-logical
Xuanwo Jul 23, 2026
ee41c49
Merge branch 'xuanwo/asof-logical' into xuanwo/asof-proto
Xuanwo Jul 23, 2026
d36fd27
fix: keep generated protobuf code in sync
Xuanwo Jul 23, 2026
73fd599
Merge remote-tracking branch 'origin/main' into xuanwo/asof-physical
Xuanwo Jul 27, 2026
43110a3
feat: broadcast ASOF right input
Xuanwo Jul 27, 2026
a15b693
Merge branch 'xuanwo/asof-physical' into xuanwo/asof-logical
Xuanwo Jul 27, 2026
88a46e1
Merge branch 'xuanwo/asof-logical' into xuanwo/asof-proto
Xuanwo Jul 27, 2026
9556ed3
fix: remove obsolete ASOF proto import
Xuanwo Jul 27, 2026
339d785
fix: account shared ASOF build buffers once
Xuanwo Jul 27, 2026
a71c19f
Merge branch 'xuanwo/asof-physical' into xuanwo/asof-logical
Xuanwo Jul 27, 2026
3a4f995
fix: align ASOF logical contracts
Xuanwo Jul 27, 2026
ce2fbc4
Merge branch 'xuanwo/asof-logical' into xuanwo/asof-proto
Xuanwo Jul 27, 2026
34552d1
perf: compare ASOF keys without scalar materialization
Xuanwo Jul 28, 2026
c967002
Merge branch 'xuanwo/asof-physical' into xuanwo/asof-logical
Xuanwo Jul 28, 2026
aaae20f
Merge branch 'xuanwo/asof-logical' into xuanwo/asof-proto
Xuanwo Jul 28, 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
51 changes: 50 additions & 1 deletion 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 @@ -1865,6 +1866,53 @@ 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 right_output_indices =
(0..join.right.schema().fields().len()).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 @@ -2364,6 +2412,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
79 changes: 75 additions & 4 deletions datafusion/expr/src/logical_plan/builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -31,10 +31,10 @@ use crate::expr_rewriter::{
rewrite_sort_cols_by_aggs,
};
use crate::logical_plan::{
Aggregate, Analyze, Distinct, DistinctOn, EmptyRelation, Explain, Filter, Join,
JoinConstraint, JoinType, Limit, LogicalPlan, Partitioning, PlanType, Prepare,
Projection, Repartition, Sort, SubqueryAlias, TableScanBuilder, Union, Unnest,
Values, Window,
Aggregate, Analyze, AsOfJoin, AsOfMatch, Distinct, DistinctOn, EmptyRelation,
Explain, Filter, Join, JoinConstraint, JoinType, Limit, LogicalPlan, Partitioning,
PlanType, Prepare, Projection, Repartition, Sort, SubqueryAlias, TableScanBuilder,
Union, Unnest, Values, Window,
};
use crate::select_expr::SelectExpr;
use crate::utils::{
Expand Down Expand Up @@ -1007,6 +1007,68 @@ impl LogicalPlanBuilder {
)
}

/// Apply a left-preserving ASOF join using equality expressions and one
/// ordered match condition.
pub fn asof_join(
self,
right: LogicalPlan,
on: Vec<(Expr, Expr)>,
match_condition: AsOfMatch,
) -> Result<Self> {
self.asof_join_with_constraint(right, on, match_condition, JoinConstraint::On)
}

/// Apply a left-preserving ASOF join using `USING` equality keys.
pub fn asof_join_using(
self,
right: LogicalPlan,
using_keys: Vec<Column>,
match_condition: AsOfMatch,
) -> Result<Self> {
let on = using_keys
.into_iter()
.map(|key| {
let left = Self::normalize(&self.plan, key.clone())?;
let right = Self::normalize(&right, key)?;
Ok((Expr::Column(left), Expr::Column(right)))
})
.collect::<Result<_>>()?;
self.asof_join_with_constraint(right, on, match_condition, JoinConstraint::Using)
}

fn asof_join_with_constraint(
self,
right: LogicalPlan,
on: Vec<(Expr, Expr)>,
match_condition: AsOfMatch,
join_constraint: JoinConstraint,
) -> Result<Self> {
let normalize = |expr, schema: &DFSchema| {
normalize_col_with_schemas_and_ambiguity_check(expr, &[&[schema]], &[])
};
let on = on
.into_iter()
.map(|(left, right_expr)| {
Ok((
normalize(left, self.plan.schema())?,
normalize(right_expr, right.schema())?,
))
})
.collect::<Result<_>>()?;
let match_condition = AsOfMatch {
left: normalize(match_condition.left, self.plan.schema())?,
op: match_condition.op,
right: normalize(match_condition.right, right.schema())?,
};
Ok(Self::new(LogicalPlan::AsOfJoin(AsOfJoin::try_new(
self.plan,
Arc::new(right),
on,
match_condition,
join_constraint,
)?)))
}

pub(crate) fn normalize(plan: &LogicalPlan, column: Column) -> Result<Column> {
if column.relation.is_some() {
// column is already normalized
Expand Down Expand Up @@ -1776,6 +1838,15 @@ pub fn build_join_schema(
dfschema.with_functional_dependencies(func_dependencies)
}

/// Creates the schema for a left-preserving ASOF join.
///
/// Both `ON` and `USING` preserve all qualified input fields. SQL wildcard
/// expansion handles the unqualified `USING` key as a single column.
pub fn build_asof_join_schema(left: &DFSchema, right: &DFSchema) -> Result<DFSchema> {
build_join_schema(left, right, &JoinType::Left)?
.with_functional_dependencies(left.functional_dependencies().clone())
}

/// (Re)qualify the sides of a join if needed, i.e. if the columns from one side would otherwise
/// conflict with the columns from the other.
/// This is especially useful for queries that come as Substrait, since Substrait doesn't currently allow specifying
Expand Down
23 changes: 19 additions & 4 deletions datafusion/expr/src/logical_plan/display.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,10 +21,10 @@ use std::collections::HashMap;
use std::fmt;

use crate::{
Aggregate, DescribeTable, Distinct, DistinctOn, DmlStatement, Expr, Filter, Join,
Limit, LogicalPlan, Partitioning, Projection, RecursiveQuery, Repartition, Sort,
Subquery, SubqueryAlias, TableProviderFilterPushDown, TableScan, Unnest, Values,
Window, expr_vec_fmt,
Aggregate, AsOfJoin, DescribeTable, Distinct, DistinctOn, DmlStatement, Expr, Filter,
Join, Limit, LogicalPlan, Partitioning, Projection, RecursiveQuery, Repartition,
Sort, Subquery, SubqueryAlias, TableProviderFilterPushDown, TableScan, Unnest,
Values, Window, expr_vec_fmt,
};

use crate::dml::CopyTo;
Expand Down Expand Up @@ -493,6 +493,21 @@ impl<'a, 'b> PgJsonVisitor<'a, 'b> {
"Filter": format!("{}", filter_expr)
})
}
LogicalPlan::AsOfJoin(AsOfJoin {
on,
match_condition,
join_constraint,
..
}) => {
let join_expr: Vec<String> =
on.iter().map(|(l, r)| format!("{l} = {r}")).collect();
json!({
"Node Type": "AsOf Join",
"Join Constraint": format!("{join_constraint:?}"),
"Join Keys": join_expr.join(", "),
"Match Condition": match_condition.to_string(),
})
}
LogicalPlan::Repartition(Repartition {
partitioning_scheme,
..
Expand Down
16 changes: 8 additions & 8 deletions datafusion/expr/src/logical_plan/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,8 +28,8 @@ pub mod tree_node;

pub use builder::{
LogicalPlanBuilder, LogicalPlanBuilderOptions, LogicalTableSource, UNNAMED_TABLE,
build_join_schema, requalify_sides_if_needed, table_scan, union,
wrap_projection_for_join_if_necessary,
build_asof_join_schema, build_join_schema, requalify_sides_if_needed, table_scan,
union, wrap_projection_for_join_if_necessary,
};
pub use ddl::{
CreateCatalog, CreateCatalogSchema, CreateExternalTable, CreateFunction,
Expand All @@ -41,12 +41,12 @@ pub use dml::{
WriteOp,
};
pub use plan::{
Aggregate, Analyze, ColumnUnnestList, DescribeTable, Distinct, DistinctOn,
EmptyRelation, Explain, ExplainOption, Extension, FetchType, Filter, Join,
JoinConstraint, JoinType, Limit, LogicalPlan, Partitioning, PlanType, Projection,
RangePartitioning, RecursiveQuery, Repartition, SkipType, Sort, StringifiedPlan,
Subquery, SubqueryAlias, TableScan, TableScanBuilder, ToStringifiedPlan, Union,
Unnest, Values, Window, projection_schema,
Aggregate, Analyze, AsOfJoin, AsOfMatch, ColumnUnnestList, DescribeTable, Distinct,
DistinctOn, EmptyRelation, Explain, ExplainOption, Extension, FetchType, Filter,
Join, JoinConstraint, JoinType, Limit, LogicalPlan, Partitioning, PlanType,
Projection, RangePartitioning, RecursiveQuery, Repartition, SkipType, Sort,
StringifiedPlan, Subquery, SubqueryAlias, TableScan, TableScanBuilder,
ToStringifiedPlan, Union, Unnest, Values, Window, projection_schema,
};
pub use statement::{
Deallocate, Execute, Prepare, ResetVariable, SetVariable, Statement,
Expand Down
Loading
Loading