Skip to content
Open
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
49 changes: 48 additions & 1 deletion datafusion/core/src/physical_planner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,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 @@ -1790,6 +1791,51 @@ 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,
)?,
);
Arc::new(AsOfJoinExec::try_new(
physical_left,
physical_right,
join_on,
match_condition,
None,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Here is something I don't fully understand — it would be great if other reviewers have the knowledge to double-check, or can point me to something that would help me understand it more easily. (Just flagging where I don't have full confidence; I don't think this is a blocker for this PR.)

Specifically, I don't follow the whole lifecycle of projection pushdown.

Let's say the downstream only requires a subset of the columns inside AsOfJoinExec. I imagine the logical optimizer should try to keep the projection list inside the LogicalPlan, so that we can directly build the physical plan node with the required projection indices.

Here, when building the initial physical plan, the projection list is None, and it seems to depend on a later physical optimizer pass to finish the work.

Perhaps there is some room to simplify that process.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good catch. I’ll handle that in a separate follow-up.

)?)
}
LogicalPlan::RecursiveQuery(RecursiveQuery {
name,
is_distinct,
Expand Down Expand Up @@ -2291,6 +2337,7 @@ fn extract_dml_filters(
| LogicalPlan::Sort(_)
| LogicalPlan::Union(_)
| LogicalPlan::Join(_)
| LogicalPlan::AsOfJoin(_)
| LogicalPlan::Repartition(_)
| LogicalPlan::Aggregate(_)
| LogicalPlan::Window(_)
Expand Down
78 changes: 74 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,
Comment on lines +1015 to +1016

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think this is the user-facing API (see the example code in the join_on comment above), so it should be easy to use.
Currently it asks callers to manually normalize the args; instead, we should validate them internally.

Perhaps:

on: Expr,
match_condition: Expr,

For example, if the SQL input is on t1.v1 < t2.v1, which is not a supported ON clause for an asof join:

  • Existing design: we have to do partial validation in the SQL->LogicalPlan binding
  • Alternative: we can consolidate the validation in one place

But I think we can proceed as is and potentially change it in the SQL integration PR. It would be obvious which design is better once we try to build a plan from SQL. And we're likely to do that before the next release, so API changes are fine.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed. Should be fixed in #23830

) -> 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,14 @@ 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)
}

/// (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