Repository navigation
feat: support col reference for percentiles - #25337
dd-annarose wants to merge 7 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.64% -0.02%
==========================================
Files 1147 1147
Lines 446357 452749 +6392
Branches 446357 452749 +6392
==========================================
+ Hits 368971 374165 +5194
- Misses 54997 55925 +928
- Partials 22389 22659 +270 ☔ 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 🤦🏻♀️) |
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