Skip to content
Merged
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
59 changes: 57 additions & 2 deletions datafusion/core/benches/sql_planner_extended.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
use arrow::array::{ArrayRef, RecordBatch};
use arrow_schema::DataType;
use arrow_schema::TimeUnit::Nanosecond;
use arrow_schema::{Field, Schema};
use criterion::{BenchmarkId, Criterion, criterion_group, criterion_main};
use datafusion::prelude::{DataFrame, SessionContext};
use datafusion_catalog::MemTable;
Expand Down Expand Up @@ -69,6 +70,28 @@ fn register_string_table(ctx: &SessionContext, num_columns: usize, num_rows: usi
ctx.register_table("t", Arc::new(table)).unwrap();
}

/// Registers a table `name` with `num_columns` Int32 columns (`a0..aN`) plus an
/// `arr` `List<Int32>` column, used to benchmark `UNNEST` planning.
fn register_list_table(ctx: &SessionContext, name: &str, num_columns: usize) {
let mut fields: Vec<Field> = (0..num_columns)
.map(|i| Field::new(format!("a{i}"), DataType::Int32, true))
.collect();
fields.push(Field::new(
"arr",
DataType::List(Arc::new(Field::new("item", DataType::Int32, true))),
true,
));
let schema = Arc::new(Schema::new(fields));
let table = MemTable::try_new(schema, vec![vec![]]).unwrap();
ctx.register_table(name, Arc::new(table)).unwrap();
}

/// Create a logical plan from the specified sql (parse + plan only, NO
/// analysis or optimization). Isolates SQL planner cost.
fn logical_plan(ctx: &SessionContext, rt: &Runtime, sql: &str) {
black_box(rt.block_on(ctx.state().create_logical_plan(sql)).unwrap());
}

/// Build a dataframe for testing logical plan optimization
fn build_test_data_frame(ctx: &SessionContext, rt: &Runtime) -> DataFrame {
register_string_table(ctx, 100, 1000);
Expand Down Expand Up @@ -387,8 +410,8 @@ fn criterion_benchmark(c: &mut Criterion) {
let case_heavy_left_join_df = build_case_heavy_left_join_df(&case_heavy_ctx, &rt);

// really slow :(
let mut group = c.benchmark_group("sample_size_5");
group.sample_size(5);
let mut group = c.benchmark_group("sample_size_10");
group.sample_size(10);
group.bench_function("logical_plan_optimize", |b| {
b.iter(|| {
let df_clone = df.clone();
Expand Down Expand Up @@ -525,6 +548,38 @@ fn criterion_benchmark(c: &mut Criterion) {
);
})
});

// UNNEST planning. The unnest rewrite builds an inner projection that is
// deduplicated on every pushed expression, so wide select lists stress it.
// See `push_projection_dedupl` in datafusion/sql/src/utils.rs.
let unnest_ctx = SessionContext::new();
register_list_table(&unnest_ctx, "t_list200", 200);
register_list_table(&unnest_ctx, "t_list1000", 1000);

// Bare columns next to an unnest: each column is pushed once.
for num_columns in [200usize, 1000] {
let cols = (0..num_columns)
.map(|i| format!("a{i}"))
.collect::<Vec<_>>()
.join(", ");
let query = format!("SELECT unnest(arr), {cols} FROM t_list{num_columns}");
c.bench_function(&format!("logical_unnest_plus_{num_columns}_columns"), |b| {
b.iter(|| logical_plan(&unnest_ctx, &rt, &query))
});
}

// Columns inside expressions that contain an unnest: exercises both the
// column path and the repeated unnest placeholder alias path.
{
let exprs = (0..200)
.map(|i| format!("unnest(arr) + a{i}"))
.collect::<Vec<_>>()
.join(", ");
let query = format!("SELECT {exprs} FROM t_list200");
c.bench_function("logical_unnest_in_200_exprs", |b| {
b.iter(|| logical_plan(&unnest_ctx, &rt, &query))
});
}
}

criterion_group!(benches, criterion_benchmark);
Expand Down
16 changes: 8 additions & 8 deletions datafusion/sql/src/select.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,10 +23,10 @@ use crate::planner::{ContextProvider, PlannerContext, SqlToRel};
use crate::query::to_order_by_exprs_with_select;
use crate::utils::{
CheckColumnsMustReferenceAggregatePurpose, CheckColumnsSatisfyExprsPurpose,
check_columns_satisfy_exprs, extract_aliases, rebase_expr, resolve_aliases_to_exprs,
resolve_columns, resolve_positions_to_exprs, rewrite_recursive_unnest_bottom_up,
rewrite_recursive_unnests_bottom_up, substitute_top_level_alias,
substitute_top_level_aliases_in_sorts,
DedupedProjection, check_columns_satisfy_exprs, extract_aliases, rebase_expr,
resolve_aliases_to_exprs, resolve_columns, resolve_positions_to_exprs,
rewrite_recursive_unnest_bottom_up, rewrite_recursive_unnests_bottom_up,
substitute_top_level_alias, substitute_top_level_aliases_in_sorts,
};

use arrow::datatypes::DataType;
Expand Down Expand Up @@ -639,7 +639,7 @@ impl<S: ContextProvider> SqlToRel<'_, S> {
let mut unnest_columns = IndexMap::new();
// from which columns used for projection, before the unnest happen
// including non unnest columns and unnest columns
let mut inner_projection_exprs = vec![];
let mut inner_projection_exprs = DedupedProjection::default();
let mut outer_expr_groups =
Vec::with_capacity(intermediate_expr_groups.len());

Expand Down Expand Up @@ -704,7 +704,7 @@ impl<S: ContextProvider> SqlToRel<'_, S> {
}

intermediate_plan = LogicalPlanBuilder::from(intermediate_plan)
.project(inner_projection_exprs)?
.project(inner_projection_exprs.into_exprs())?
.unnest_columns_with_options(unnest_col_vec, unnest_options)?
.build()?;
intermediate_expr_groups = outer_expr_groups;
Expand Down Expand Up @@ -807,7 +807,7 @@ impl<S: ContextProvider> SqlToRel<'_, S> {

loop {
let mut unnest_columns = IndexMap::new();
let mut inner_projection_exprs = vec![];
let mut inner_projection_exprs = DedupedProjection::default();

let outer_projection_exprs = rewrite_recursive_unnests_bottom_up(
&intermediate_plan,
Expand Down Expand Up @@ -842,7 +842,7 @@ impl<S: ContextProvider> SqlToRel<'_, S> {
columns
}
};
projection_exprs.extend(inner_projection_exprs);
projection_exprs.extend(inner_projection_exprs.into_exprs());

let mut unnest_col_vec = vec![];

Expand Down
Loading
Loading