Repository navigation
feat: support col reference for percentiles - #25337
dd-annarose wants to merge 9 commits into
Conversation
|
Thank you for opening this pull request! Reviewer note: cargo-semver-checks reported the current version number is not SemVer-compatible with the changes in this pull request (compared against the base branch). Details |
Codecov Report❌ Patch coverage is Additional details and impacted files@@ Coverage Diff @@
## main #25337 +/- ##
==========================================
+ Coverage 82.66% 82.78% +0.11%
==========================================
Files 1147 1147
Lines 446357 451163 +4806
Branches 446357 451163 +4806
==========================================
+ Hits 368971 373481 +4510
+ Misses 54997 54942 -55
- Partials 22389 22740 +351 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
alamb
left a comment
There was a problem hiding this comment.
Thanks @dd-annarose
Can we do this at planning time instead? that would likely be must faster as it wouldn't have to inspect the inputs.
In the example you showed
SELECT
approx_percentile_cont(y, m) AS median
FROM (
SELECT t.x + 1 as y, 0.5 as m
FROM (
VALUES (10)
) AS t(x)
)I think the subquery could be flattened to
SELECT
approx_percentile_cont(t.x + 1, 0.5) AS median
FROM
VALUES (10)( in fact I am surprised it isn't)
|
@alamb We could, but we wouldn't support queries like this one (supported by Trino for example): I don't have a Trino parity agenda, but this can be used more broadly when wanting to use a "real" column reference as percentile. Happy to discuss! |
|
Also: this is supported by spark |
asolimando
left a comment
There was a problem hiding this comment.
@dd-annarose thanks for working on this, I think it's good to support expressions beyond literals, as Postgres, Spark and Trino do. I left two comments, one minor potential perf improvement, and one blocking issue/error, the rest LGTM!
147bbb8 to
e778c2f
Compare
Adds a PercentileParam and PercentileParamState to accept column references or projections for the percentile argument. # Conflicts: # datafusion/functions-aggregate/src/percentile_cont.rs
replace non-deterministic metrics with <slt:ignore>
|
run benchmark clickbench_partitioned |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing annarose/approx_percentile (44a9b2b) to d040501 (merge-base) diff Run configurationrun benchmark clickbench_partitionedResults will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing annarose/approx_percentile (44a9b2b) to d040501 (merge-base) diff Run configurationrun benchmark clickbench_partitionedCPU Details (lscpu)Details
Resource Usageclickbench_partitioned — base (merge-base)
clickbench_partitioned — branch
File an issue against this benchmark runner |
cleaner code
0aa3562 to
8088f8c
Compare
Regarding this, I'm seeing that this is actually not supported by Spark. Did a quick test with a local Spark container running 4.0.1 and this is what I'm seeing:
In the doc you linked https://spark.apache.org/docs/latest/api/python/reference/pyspark.sql/api/pyspark.sql.functions.approx_percentile.html#pyspark-sql-functions-approx-percentile, it looks like that For reference here is the Spark code that rejects these kind of queries in case of a NON_FOLDABLE_INPUT. Given that this is supported by Trino, and that it's still API compatible with Spark and Postgres, I'd still be in favor of merging this, however, I'm not comfortable just pulling this in without an extra pair of eyes from a PMC. cc @alamb as you reviewed this PR before. |
|
(sorry I closed the wrong PR 🤦🏻♀️) |
gabotechs
left a comment
There was a problem hiding this comment.
Leaving here a +1, as on my side there's no more blockers here and I think it's worth pulling this in.
I'll still give a chance for other contributors to chime in.
| return Ok(()); | ||
| } | ||
|
|
||
| let agg_fn_name = self.aggregate_fn_name.clone(); |
There was a problem hiding this comment.
This cloning could be avoided. Just pass the &self.aggregate_fn_name where needed.
| match array.data_type() { | ||
| DataType::Float64 | DataType::Float32 => { | ||
| let array = cast(array, &DataType::Float64)?; | ||
| Ok(downcast_value!(array, Float64Array).clone()) |
There was a problem hiding this comment.
Is .clone() really needed here ?
There was a problem hiding this comment.
operator cannot convert from
&PrimitiveArray<_>toPrimitiveArray<_>
downcast_value!(array, Float64Array) returns a &Florat64Array
There was a problem hiding this comment.
I could use to_owned but it's essentially cloning
|
|
||
| fn update_batch(&mut self, values: &[ArrayRef]) -> Result<()> { | ||
| self.approx_percentile_cont_accumulator | ||
| .resolve_percentile(&values[2])?; |
There was a problem hiding this comment.
What guarantees that the passed values has at least three elements ?
There was a problem hiding this comment.
values contains the arguments to this aggregate function.
The aggregate function's accumulator() also returns an error if there are not 3 or 4 arguments (approx_percentile_cont_with_weight requires three or four arguments: value, weight, percentile[, centroids]).
But I can add an assertion here if you'd prefer.
|
|
||
| /// Percentile argument for `aggregate_fn_name` and its state. | ||
| #[derive(Debug)] | ||
| pub struct PercentileParam { |
There was a problem hiding this comment.
The API is strange - a pub struct with pub(crate) fields.
Since the struct is not exported in lib.rs it should be pub(crate).
There was a problem hiding this comment.
I decided to keep it as pub and make PercentileParam::try_new pub too in order to preserve API contract for ApproxPercentileAccumulator's constructor.
| ) | ||
| } | ||
| data_type => { | ||
| return plan_err!( |
There was a problem hiding this comment.
resolve() is called at execution time, no ? This should be exec_err!()
There was a problem hiding this comment.
ah good catch, thanks!
| internal_datafusion_err!("expected a non-null percentile value") | ||
| })?; | ||
| if batch_min != batch_max { | ||
| return plan_err!( |
There was a problem hiding this comment.
| return plan_err!( | |
| return exec_err!( |
| )?), | ||
| is_desc, | ||
| }), | ||
| Err(_) => Ok(PercentileParam { |
There was a problem hiding this comment.
| Err(_) => Ok(PercentileParam { | |
| Err(_) if !datafusion_physical_expr::utils::collect_columns(expr).is_empty() => Ok(PercentileParam { |
PercentileParamState::Pending state should be used only for columns, right ?
| })?; | ||
| let batch_max = batch_max.ok_or_else(|| { | ||
| internal_datafusion_err!("expected a non-null percentile value") | ||
| })?; |
There was a problem hiding this comment.
| })?; | |
| })?; | |
| if batch_min.is_nan() || batch_max.is_nan() { | |
| return exec_err!( | |
| "Percentile value must be between 0.0 and 1.0 inclusive, NaN is invalid" | |
| ); | |
| } |
It would be good to add some tests with NaN and infinite floats just to make sure that they are properly handled.
| 02)--AggregateExec: mode=Final, gby=[], aggr=[count(Int64(1)), sum(total)], metrics=[output_rows=1, elapsed_compute=<slt:ignore>, output_bytes=16.0 B, output_batches=1, agg_expr_0_arguments_time=<slt:ignore>, agg_expr_0_evaluate_time=<slt:ignore>, agg_expr_0_merge_time=<slt:ignore>, agg_expr_1_arguments_time=<slt:ignore>, agg_expr_1_evaluate_time=<slt:ignore>, agg_expr_1_merge_time=<slt:ignore>] | ||
| 03)----CoalescePartitionsExec, metrics=[output_rows=4, elapsed_compute=<slt:ignore>, output_bytes=64.0 B, output_batches=4] | ||
| 04)------AggregateExec: mode=Partial, gby=[], aggr=[count(Int64(1)), sum(total)], metrics=[output_rows=4, elapsed_compute=<slt:ignore>, output_bytes=64.0 B, output_batches=4, agg_expr_0_arguments_time=<slt:ignore>, agg_expr_0_state_time=<slt:ignore>, agg_expr_0_update_time=<slt:ignore>, agg_expr_1_arguments_time=<slt:ignore>, agg_expr_1_state_time=<slt:ignore>, agg_expr_1_update_time=<slt:ignore>] | ||
| 05)--------ProjectionExec: expr=[sum(t.v)@1 as total], metrics=[output_rows=100.0 K, elapsed_compute=<slt:ignore>, output_bytes=787.4 KB, output_batches=788, expr_0_eval_time=<slt:ignore>] |
There was a problem hiding this comment.
Are output_bytes=787.4 KB and output_batches=788 important for the test here ?
If not then I'd suggest to <slt:ignore> them.
added sqllogictests for inf and NaN
Which issue does this PR close?
percentile_contandapprox_percentile_cont#24337.Rationale for this change
When using a projection as the percentile argument to
approx_percentile_contorpercentile_cont, planning fails because only literals are supported.It makes sense to accept projections that are constant across all batches but not labelled as literals (although they practically are).
For example, this query cannot be planned, even though the percentile is constant:
What changes are included in this PR?
Introduce
PercentileParamandPercentileParamStateto handle column references or projections for the percentile argument.Resolution of
percentilecannot always be done: for example, on an empty record batch, we won't be able to get the value from the columns.PercentileParamStatemakes sure the percentile resolution is both flexible enough to wait to be able to resolve the parameter AND applies the strict constraints that are required (constantFloat32orFloat64).The percentile is also included in state fields so that the merge_batch in the final Accumulator can process without any issues.
What is the testing strategy for this PR?
This new feature is covered by the
sqllogictestcases added inaggregate.slt.Unit tests have also been added to
approx_percentile_cont.rsandpercentile_cont.rs.Are there any user-facing changes?
No breaking changes.
Benchark: ClickBench
Memory profiling for clickbench_partitioned
this branch
main