-
Notifications
You must be signed in to change notification settings - Fork 2k
DataFrame API: allow aggregate functions in select() (#17874) #21021
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -57,6 +57,7 @@ use datafusion_common::{ | |
| plan_datafusion_err, plan_err, unqualified_field_not_found, | ||
| }; | ||
| use datafusion_expr::select_expr::SelectExpr; | ||
| use datafusion_expr::utils::find_aggregate_exprs; | ||
| use datafusion_expr::{ | ||
| ExplainOption, SortExpr, TableProviderFilterPushDown, UNNAMED_TABLE, case, | ||
| dml::InsertOp, | ||
|
|
@@ -410,21 +411,76 @@ impl DataFrame { | |
| expr_list: impl IntoIterator<Item = impl Into<SelectExpr>>, | ||
| ) -> Result<DataFrame> { | ||
| let expr_list: Vec<SelectExpr> = | ||
| expr_list.into_iter().map(|e| e.into()).collect::<Vec<_>>(); | ||
| expr_list.into_iter().map(|e| e.into()).collect(); | ||
|
|
||
| // Extract plain expressions | ||
| let expressions = expr_list.iter().filter_map(|e| match e { | ||
| SelectExpr::Expression(expr) => Some(expr), | ||
| _ => None, | ||
| }); | ||
|
|
||
| let window_func_exprs = find_window_exprs(expressions); | ||
| let plan = if window_func_exprs.is_empty() { | ||
| // Apply window functions first | ||
| let window_func_exprs = find_window_exprs(expressions.clone()); | ||
|
|
||
| let mut plan = if window_func_exprs.is_empty() { | ||
| self.plan | ||
| } else { | ||
| LogicalPlanBuilder::window_plan(self.plan, window_func_exprs)? | ||
| }; | ||
|
|
||
| let project_plan = LogicalPlanBuilder::from(plan).project(expr_list)?.build()?; | ||
| // Collect aggregate expressions | ||
| let aggr_exprs = find_aggregate_exprs(expressions.clone()); | ||
|
|
||
| // Check if any expression is non-aggregate | ||
| let has_non_aggregate_expr = expressions | ||
| .clone() | ||
| .any(|expr| find_aggregate_exprs(std::iter::once(expr)).is_empty()); | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. What about aggregate expr + non-aggregate one ? let res = df.select(vec![
count(col("c9")).alias("count_c9") + lit(1)
])?;I'd expect 101 but it returns 100 |
||
|
|
||
| // Fallback to projection: | ||
| // - already aggregated | ||
| // - contains non-aggregate expressions | ||
| // - no aggregates at all | ||
| if matches!(plan, LogicalPlan::Aggregate(_)) | ||
| || has_non_aggregate_expr | ||
| || aggr_exprs.is_empty() | ||
| { | ||
| let project_plan = | ||
| LogicalPlanBuilder::from(plan).project(expr_list)?.build()?; | ||
|
|
||
| return Ok(DataFrame { | ||
| session_state: self.session_state, | ||
| plan: project_plan, | ||
| projection_requires_validation: false, | ||
| }); | ||
| } | ||
|
|
||
| // Build Aggregate node | ||
| let aggr_exprs: Vec<Expr> = aggr_exprs | ||
| .into_iter() | ||
| .enumerate() | ||
| .map(|(i, expr)| expr.alias(format!("__agg_{i}"))) | ||
| .collect(); | ||
|
|
||
| plan = LogicalPlanBuilder::from(plan) | ||
| .aggregate(Vec::<Expr>::new(), aggr_exprs)? | ||
| .build()?; | ||
|
|
||
| // Replace aggregates with their aliases | ||
| let mut rewritten_exprs = Vec::with_capacity(expr_list.len()); | ||
| for (i, select_expr) in expr_list.into_iter().enumerate() { | ||
| match select_expr { | ||
| SelectExpr::Expression(expr) => { | ||
| let column = Expr::Column(Column::from_name(format!("__agg_{i}"))); | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
|
||
| let alias = expr.name_for_alias()?; | ||
| rewritten_exprs.push(SelectExpr::Expression(column.alias(alias))); | ||
| } | ||
| other => rewritten_exprs.push(other), | ||
| } | ||
| } | ||
|
|
||
| let project_plan = LogicalPlanBuilder::from(plan) | ||
| .project(rewritten_exprs)? | ||
| .build()?; | ||
|
|
||
| Ok(DataFrame { | ||
| session_state: self.session_state, | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -6854,3 +6854,26 @@ async fn test_duplicate_state_fields_for_dfschema_construct() -> Result<()> { | |
|
|
||
| Ok(()) | ||
| } | ||
|
|
||
| #[tokio::test] | ||
| async fn test_dataframe_api_aggregate_fn_in_select() -> Result<()> { | ||
| let df = test_table().await?; | ||
|
|
||
| let res = df.select(vec![ | ||
| count(col("c9")).alias("count_c9"), | ||
| count(cast(col("c9"), DataType::Utf8View)).alias("count_c9_str"), | ||
| ])?; | ||
|
|
||
| assert_batches_eq!( | ||
| &[ | ||
| "+----------+--------------+", | ||
| "| count_c9 | count_c9_str |", | ||
| "+----------+--------------+", | ||
| "| 100 | 100 |", | ||
| "+----------+--------------+", | ||
| ], | ||
| &res.collect().await? | ||
| ); | ||
|
|
||
| Ok(()) | ||
| } | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Please add tests for some more complex queries, e.g. |
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
find_aggregate_exprs()deduplicates the expressions.Test like:
fails with:
__agg_1is lost