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
2 changes: 1 addition & 1 deletion datafusion-testing
Submodule datafusion-testing updated 88 files
+3 −3 data/sqlite/random/expr/slt_good_0.slt
+10 −10 data/sqlite/random/expr/slt_good_10.slt
+5 −3 data/sqlite/random/expr/slt_good_100.slt
+3 −3 data/sqlite/random/expr/slt_good_101.slt
+12 −12 data/sqlite/random/expr/slt_good_102.slt
+5 −3 data/sqlite/random/expr/slt_good_103.slt
+7 −7 data/sqlite/random/expr/slt_good_104.slt
+9 −9 data/sqlite/random/expr/slt_good_105.slt
+6 −6 data/sqlite/random/expr/slt_good_106.slt
+3 −3 data/sqlite/random/expr/slt_good_107.slt
+6 −6 data/sqlite/random/expr/slt_good_108.slt
+5 −3 data/sqlite/random/expr/slt_good_109.slt
+3 −3 data/sqlite/random/expr/slt_good_11.slt
+6 −6 data/sqlite/random/expr/slt_good_110.slt
+6 −6 data/sqlite/random/expr/slt_good_111.slt
+7 −7 data/sqlite/random/expr/slt_good_112.slt
+6 −6 data/sqlite/random/expr/slt_good_113.slt
+3 −3 data/sqlite/random/expr/slt_good_116.slt
+3 −3 data/sqlite/random/expr/slt_good_118.slt
+4 −4 data/sqlite/random/expr/slt_good_119.slt
+9 −7 data/sqlite/random/expr/slt_good_12.slt
+14 −12 data/sqlite/random/expr/slt_good_13.slt
+6 −6 data/sqlite/random/expr/slt_good_14.slt
+6 −6 data/sqlite/random/expr/slt_good_16.slt
+3 −3 data/sqlite/random/expr/slt_good_17.slt
+3 −3 data/sqlite/random/expr/slt_good_2.slt
+6 −6 data/sqlite/random/expr/slt_good_22.slt
+10 −6 data/sqlite/random/expr/slt_good_23.slt
+3 −3 data/sqlite/random/expr/slt_good_24.slt
+12 −13 data/sqlite/random/expr/slt_good_25.slt
+6 −6 data/sqlite/random/expr/slt_good_26.slt
+8 −6 data/sqlite/random/expr/slt_good_27.slt
+3 −3 data/sqlite/random/expr/slt_good_28.slt
+14 −12 data/sqlite/random/expr/slt_good_29.slt
+3 −3 data/sqlite/random/expr/slt_good_3.slt
+3 −3 data/sqlite/random/expr/slt_good_30.slt
+9 −10 data/sqlite/random/expr/slt_good_31.slt
+9 −9 data/sqlite/random/expr/slt_good_33.slt
+3 −3 data/sqlite/random/expr/slt_good_34.slt
+3 −3 data/sqlite/random/expr/slt_good_36.slt
+3 −3 data/sqlite/random/expr/slt_good_38.slt
+3 −3 data/sqlite/random/expr/slt_good_39.slt
+3 −3 data/sqlite/random/expr/slt_good_4.slt
+3 −3 data/sqlite/random/expr/slt_good_41.slt
+6 −6 data/sqlite/random/expr/slt_good_44.slt
+3 −3 data/sqlite/random/expr/slt_good_45.slt
+3 −3 data/sqlite/random/expr/slt_good_47.slt
+3 −3 data/sqlite/random/expr/slt_good_49.slt
+3 −3 data/sqlite/random/expr/slt_good_5.slt
+3 −3 data/sqlite/random/expr/slt_good_50.slt
+5 −3 data/sqlite/random/expr/slt_good_53.slt
+3 −3 data/sqlite/random/expr/slt_good_55.slt
+9 −9 data/sqlite/random/expr/slt_good_56.slt
+5 −3 data/sqlite/random/expr/slt_good_57.slt
+4 −4 data/sqlite/random/expr/slt_good_58.slt
+9 −9 data/sqlite/random/expr/slt_good_6.slt
+3 −3 data/sqlite/random/expr/slt_good_60.slt
+3 −3 data/sqlite/random/expr/slt_good_61.slt
+3 −3 data/sqlite/random/expr/slt_good_62.slt
+6 −6 data/sqlite/random/expr/slt_good_63.slt
+3 −3 data/sqlite/random/expr/slt_good_64.slt
+3 −3 data/sqlite/random/expr/slt_good_66.slt
+14 −12 data/sqlite/random/expr/slt_good_67.slt
+5 −3 data/sqlite/random/expr/slt_good_68.slt
+3 −3 data/sqlite/random/expr/slt_good_69.slt
+3 −3 data/sqlite/random/expr/slt_good_7.slt
+5 −3 data/sqlite/random/expr/slt_good_72.slt
+3 −3 data/sqlite/random/expr/slt_good_75.slt
+5 −3 data/sqlite/random/expr/slt_good_77.slt
+3 −3 data/sqlite/random/expr/slt_good_78.slt
+3 −3 data/sqlite/random/expr/slt_good_79.slt
+3 −3 data/sqlite/random/expr/slt_good_8.slt
+3 −3 data/sqlite/random/expr/slt_good_80.slt
+10 −10 data/sqlite/random/expr/slt_good_81.slt
+10 −10 data/sqlite/random/expr/slt_good_82.slt
+6 −6 data/sqlite/random/expr/slt_good_83.slt
+11 −12 data/sqlite/random/expr/slt_good_84.slt
+7 −7 data/sqlite/random/expr/slt_good_85.slt
+3 −3 data/sqlite/random/expr/slt_good_88.slt
+3 −3 data/sqlite/random/expr/slt_good_89.slt
+11 −12 data/sqlite/random/expr/slt_good_92.slt
+8 −6 data/sqlite/random/expr/slt_good_93.slt
+3 −3 data/sqlite/random/expr/slt_good_94.slt
+3 −3 data/sqlite/random/expr/slt_good_96.slt
+8 −6 data/sqlite/random/expr/slt_good_97.slt
+2 −2 data/sqlite/random/groupby/slt_good_12.slt
+4 −4 data/sqlite/random/groupby/slt_good_4.slt
+5 −5 data/sqlite/random/groupby/slt_good_7.slt
6 changes: 6 additions & 0 deletions datafusion/common/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1077,6 +1077,12 @@ config_namespace! {
/// one range filter.
pub enable_piecewise_merge_join: bool, default = false

/// When set to true, an as-of join (`ASOF JOIN`) requires its inputs to be sorted by
/// the equality keys and then the as-of key. The planner adds a sort only when an input
/// is not already so ordered, letting pre-sorted inputs skip sorting entirely. When
/// false, the operator instead collects and sorts each input in memory.
pub asof_join_use_sorted_input: bool, default = true

/// The maximum estimated size in bytes for one input side of a HashJoin
/// will be collected into a single partition
pub hash_join_single_partition_threshold: usize, default = 1024 * 1024
Expand Down
4 changes: 4 additions & 0 deletions datafusion/core/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -204,6 +204,10 @@ name = "csv_load"
harness = false
name = "distinct_query_sql"

[[bench]]
harness = false
name = "asof_join_sql"

[[bench]]
harness = false
name = "push_down_filter"
Expand Down
149 changes: 149 additions & 0 deletions datafusion/core/benches/asof_join_sql.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,149 @@
// Licensed to the Apache Software Foundation (ASF) under one
// or more contributor license agreements. See the NOTICE file
// distributed with this work for additional information
// regarding copyright ownership. The ASF licenses this file
// to you under the Apache License, Version 2.0 (the
// "License"); you may not use this file except in compliance
// with the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.

//! Criterion benchmark comparing an as-of join expressed as SQL across builds.
//!
//! The query
//!
//! ```sql
//! SELECT ... FROM left l JOIN right r ON l.sym = r.sym AND l.t >= r.t
//! ```
//!
//! is run through a full `SessionContext`. On a build that carries the AsOf
//! join optimizer rule (`try_asof_join` inside `JoinSelection`), the physical
//! plan is rewritten to `AsOfJoinExec`, which emits a single nearest match per
//! left row. On a build without the rule (e.g. `branch-53`), the planner falls
//! back to a hash join plus an inequality filter, which materializes *every*
//! right row with `r.t <= l.t`.
//!
//! The two strategies are therefore "similar SQL" rather than identical
//! semantics, but running this same binary on both branches gives an
//! apples-to-apples comparison of how each executes the query. On startup the
//! benchmark prints the physical plan so the active strategy is visible in the
//! output.

use std::hint::black_box;
use std::sync::Arc;

use arrow::array::Int64Array;
use arrow::datatypes::{DataType, Field, Schema, SchemaRef};
use arrow::record_batch::RecordBatch;
use criterion::{Criterion, criterion_group, criterion_main};
use datafusion::datasource::MemTable;
use datafusion::error::Result;
use datafusion::execution::context::SessionContext;
use tokio::runtime::Runtime;

/// Left/right table sizes and partition count. Kept modest because the
/// no-asof (hash join + filter) plan materializes an inequality join whose
/// output grows with the per-`sym` partition size.
const LEFT_ROWS: usize = 50_000;
const RIGHT_ROWS: usize = 50_000;
const NUM_SYMS: usize = 2_000;

fn schema() -> SchemaRef {
Arc::new(Schema::new(vec![
Field::new("sym", DataType::Int64, false),
Field::new("t", DataType::Int64, false),
Field::new("v", DataType::Int64, false),
]))
}

fn batches(num_rows: usize, num_syms: usize, schema: &SchemaRef) -> Vec<RecordBatch> {
let syms: Vec<i64> = (0..num_rows).map(|i| (i % num_syms) as i64).collect();
let ts: Vec<i64> = (0..num_rows).map(|i| i as i64).collect();
let vals: Vec<i64> = (0..num_rows).map(|i| i as i64).collect();

let batch = RecordBatch::try_new(
Arc::clone(schema),
vec![
Arc::new(Int64Array::from(syms)),
Arc::new(Int64Array::from(ts)),
Arc::new(Int64Array::from(vals)),
],
)
.unwrap();

let batch_size = 8192;
let mut out = Vec::new();
let mut offset = 0;
while offset < batch.num_rows() {
let len = (batch.num_rows() - offset).min(batch_size);
out.push(batch.slice(offset, len));
offset += len;
}
out
}

fn make_ctx() -> Result<SessionContext> {
let ctx = SessionContext::new();
let s = schema();
let left = MemTable::try_new(Arc::clone(&s), vec![batches(LEFT_ROWS, NUM_SYMS, &s)])?;
let right =
MemTable::try_new(Arc::clone(&s), vec![batches(RIGHT_ROWS, NUM_SYMS, &s)])?;
ctx.register_table("left_tbl", Arc::new(left))?;
ctx.register_table("right_tbl", Arc::new(right))?;
Ok(ctx)
}

/// Backward as-of: for each left row, the right row(s) with `r.t <= l.t`.
///
/// All six columns are selected so no projection is pushed into the join; this
/// keeps the comparison on the raw join/match cost and sidesteps the operator's
/// (currently broken) projection path.
const SQL_INNER: &str = "SELECT l.sym, l.t, l.v, r.sym AS r_sym, r.t AS r_t, r.v AS r_v \
FROM left_tbl l JOIN right_tbl r ON l.sym = r.sym AND l.t >= r.t";
const SQL_LEFT: &str = "SELECT l.sym, l.t, l.v, r.sym AS r_sym, r.t AS r_t, r.v AS r_v \
FROM left_tbl l LEFT JOIN right_tbl r ON l.sym = r.sym AND l.t >= r.t";

fn run_sql(ctx: &SessionContext, rt: &Runtime, sql: &str) {
let df = rt.block_on(ctx.sql(sql)).unwrap();
black_box(rt.block_on(df.collect()).unwrap());
}

/// Print the physical plan once so the output records which execution strategy
/// (AsOfJoinExec vs hash join + filter) this build selected.
fn report_plan(ctx: &SessionContext, rt: &Runtime, label: &str, sql: &str) {
let plan = rt
.block_on(async {
let df = ctx.sql(sql).await?;
df.create_physical_plan().await
})
.unwrap();
let displayable = datafusion::physical_plan::displayable(plan.as_ref()).indent(true);
eprintln!("\n===== physical plan [{label}] =====\n{displayable}");
}

fn bench_asof_sql(c: &mut Criterion) {
let rt = Runtime::new().unwrap();
let ctx = make_ctx().unwrap();

report_plan(&ctx, &rt, "inner_backward", SQL_INNER);
report_plan(&ctx, &rt, "left_backward", SQL_LEFT);

let mut group = c.benchmark_group("asof_join_sql");

group.bench_function("inner_backward", |b| {
b.iter(|| run_sql(&ctx, &rt, SQL_INNER))
});
group.bench_function("left_backward", |b| b.iter(|| run_sql(&ctx, &rt, SQL_LEFT)));

group.finish();
}

criterion_group!(benches, bench_asof_sql);
criterion_main!(benches);
169 changes: 167 additions & 2 deletions datafusion/core/src/physical_planner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,8 @@ use crate::error::{DataFusionError, Result};
use crate::execution::context::{ExecutionProps, SessionState};
use crate::logical_expr::utils::generate_sort_key;
use crate::logical_expr::{
Aggregate, EmptyRelation, Join, Projection, Sort, TableScan, Unnest, Values, Window,
Aggregate, AsOfJoin, EmptyRelation, Join, Projection, Sort, TableScan, Unnest,
Values, Window,
};
use crate::logical_expr::{
Expr, LogicalPlan, Partitioning as LogicalPartitioning, PlanType, Repartition,
Expand All @@ -42,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,
AsOfJoinCondition, AsOfJoinExec, CrossJoinExec, HashJoinExec, NestedLoopJoinExec,
PartitionMode, SortMergeJoinExec,
};
use crate::physical_plan::limit::{GlobalLimitExec, LocalLimitExec};
use crate::physical_plan::projection::{ProjectionExec, ProjectionExpr};
Expand Down Expand Up @@ -1709,6 +1711,137 @@ impl DefaultPhysicalPlanner {
join
}
}
LogicalPlan::AsOfJoin(AsOfJoin {
left,
right,
on,
match_condition,
join_type,
null_equality,
..
}) => {
let [physical_left, physical_right] = children.two()?;

let left_df_schema = left.schema();
let right_df_schema = right.schema();
let execution_props = session_state.execution_props();

// Equijoin keys: `.0` references the LEFT input, `.1` the RIGHT.
let join_on = on
.iter()
.map(|(l, r)| {
let l = create_physical_expr(l, left_df_schema, execution_props)?;
let r =
create_physical_expr(r, right_df_schema, execution_props)?;
Ok((l, r))
})
.collect::<Result<join_utils::JoinOn>>()?;

// The match condition must be a single inequality comparison.
let Expr::BinaryExpr(BinaryExpr {
left: match_left,
op,
right: match_right,
}) = match_condition
else {
return plan_err!(
"AsOfJoin match_condition must be a binary inequality expression, got {match_condition}"
);
};

let mut op = *op;
if !matches!(
op,
Operator::Lt | Operator::LtEq | Operator::Gt | Operator::GtEq
) {
return plan_err!(
"AsOfJoin match_condition requires an inequality operator (<, <=, >, >=), got {op:?}"
);
}

fn reverse_ineq(op: Operator) -> Operator {
match op {
Operator::Lt => Operator::Gt,
Operator::LtEq => Operator::GtEq,
Operator::Gt => Operator::Lt,
Operator::GtEq => Operator::LtEq,
_ => op,
}
}

#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum Side {
Left,
Right,
Both,
}

let side_of = |e: &Expr| -> Result<Side> {
let cols = e.column_refs();
let any_left = cols
.iter()
.any(|c| left_df_schema.index_of_column(c).is_ok());
let any_right = cols
.iter()
.any(|c| right_df_schema.index_of_column(c).is_ok());
Ok(match (any_left, any_right) {
(true, false) => Side::Left,
(false, true) => Side::Right,
(true, true) => Side::Both,
(false, false) => {
return plan_err!(
"AsOfJoin match_condition operand must reference an input column"
);
}
})
};

// Orient the condition so it reads `left_input op right_input`,
// flipping the operator if it was written the other way around.
let mut lhs_logical = match_left.as_ref();
let mut rhs_logical = match_right.as_ref();
let lhs_side = side_of(lhs_logical)?;
let rhs_side = side_of(rhs_logical)?;

if lhs_side == Side::Both || rhs_side == Side::Both {
return plan_err!(
"AsOfJoin match_condition operands must each reference a single input side"
);
}

if lhs_side == Side::Right && rhs_side == Side::Left {
std::mem::swap(&mut lhs_logical, &mut rhs_logical);
op = reverse_ineq(op);
} else if !(lhs_side == Side::Left && rhs_side == Side::Right) {
return plan_err!(
"AsOfJoin match_condition must compare the left input against the right input"
);
}

let condition_left =
create_physical_expr(lhs_logical, left_df_schema, execution_props)?;
let condition_right =
create_physical_expr(rhs_logical, right_df_schema, execution_props)?;
let asof_condition =
AsOfJoinCondition::try_new(condition_left, op, condition_right)?;

let use_sorted_input = session_state
.config_options()
.optimizer
.asof_join_use_sorted_input;
Arc::new(
AsOfJoinExec::try_new(
physical_left,
physical_right,
join_on,
asof_condition,
*join_type,
None,
*null_equality,
)?
.with_required_input_ordering(use_sorted_input),
)
}
LogicalPlan::RecursiveQuery(RecursiveQuery {
name, is_distinct, ..
}) => {
Expand Down Expand Up @@ -2167,6 +2300,7 @@ fn extract_dml_filters(
| LogicalPlan::Sort(_)
| LogicalPlan::Union(_)
| LogicalPlan::Join(_)
| LogicalPlan::AsOfJoin(_)
| LogicalPlan::Repartition(_)
| LogicalPlan::Aggregate(_)
| LogicalPlan::Window(_)
Expand Down Expand Up @@ -3784,6 +3918,37 @@ mod tests {
Ok(())
}

#[tokio::test]
async fn test_asof_join_planning() -> Result<()> {
let schema = Schema::new(vec![
Field::new("k", DataType::Int32, false),
Field::new("t", DataType::Int64, false),
]);

let left = scan_empty(Some("l"), &schema, None)?.build()?;
let right = scan_empty(Some("r"), &schema, None)?.build()?;

let asof = AsOfJoin::try_new(
Arc::new(left),
Arc::new(right),
vec![(col("l.k"), col("r.k"))],
col("l.t").gt_eq(col("r.t")),
JoinType::Inner,
datafusion_common::NullEquality::NullEqualsNothing,
)?;
let logical_plan = LogicalPlan::AsOfJoin(asof);

let session_state = make_session_state();
let planner = DefaultPhysicalPlanner::default();
let physical_plan = planner
.create_physical_plan(&logical_plan, &session_state)
.await?;

let displayed = displayable(physical_plan.as_ref()).indent(true).to_string();
assert_contains!(displayed, "AsOfJoinExec");
Ok(())
}

#[tokio::test]
async fn test_explain() {
let schema = Schema::new(vec![Field::new("id", DataType::Int32, false)]);
Expand Down
Loading
Loading