Repository navigation
Introduce a way to represent constrained statistics / bounds on values in Statistics #8078
Description
Activity
- changed the title
[-]Introduce a way to represent constrained staitistics[/-][+]Introduce a way to represent constrained statistics / bounds on values in Statistics[/+]on Nov 7, 2023 A further wrinkle that may be worth thinking about, even if the solution is not to do anything, is that some formats, including parquet, special case certain floating point values. There is some discussion on apache/parquet-format#196 to try to make the statistics actually usable, and one potential outcome might be to just use totalOrder, so it may even be no action is required, but just an FYI
Reacted by Andrew LambI have two proposals:
If we have the cases where the lower bound of a range might be inexact (or perhaps absent) and the upper bound is exact—or vice versa—we could define the following:
enum Estimation { Exact(T), Inexact(T), Absent, } enum Precision { Estimation, // min - max Range(Estimation, Estimation) }- If we don't need to consider the exactness of bounds and a simple interval suffices, this alternative could be used:
enum Precision { Exact(T), Inexact(T), Absent, Between(Interval) }enum Precision { Exact(T), Inexact(T), Absent, Between(Interval) }@berkaysynnada under what circumstances would we use
Inexact? All the cases I can think ofInexactnow would be handled byBetween(potentially with unbounded intervals)Actually, upon more thought I think the best next step here is for me to prototype what this would look like. I'll do that and report back
enum Precision { Exact(T), Inexact(T), Absent, Between(Interval) }@berkaysynnada under what circumstances would we use
Inexact? All the cases I can think ofInexactnow would be handled byBetween(potentially with unbounded intervals)You are right, I cannot find any case to use Inexact. Any exact value can be relaxed to positive or negative side using intervals.
In proposal two, you still need
Inexact; e.g. cases involving selectivity. Assume that you have a columnxwhose bounds are[1, 100]. After applying a filter ofx <= 10, the row count estimate is an inexact estimate of0.1 * NwhereNwas the estimate prior to the filter.On the contrary: If you have
Between, you do not needExactanymore -- AnExactvalue is simply aBetweenvalue with a singleton interval inside.BTW @berkaysynnada's first proposal is the more general one: It allows range estimations where either bound can be exact or inexact estimations. His second proposal is slightly simpler to implement (and leverages the interval library), but it constrains range estimations to exact bounds.
When we decide to support slightly more complicated filters, say
x <= y, we will probably need range estimations where both sides could be inexact. Therefore, I think it makes sense to keep things flexible and go with the first proposal. Obv. happy to discuss other perspectives.I think there is a mismatch between the current use of
Precision::Inexactand its documentationSpecifically,
Precision::Inexactappears to be treated as a "conservative estimate" in several places (which is actually what we need in IOx), so perhaps another option would be to renamePrecision::ExacttoPrecision::Conservativeand document that it is a conservative estimate of the actual value 🤔For example the comments on
Precision::InexactsayHowever, it is used to skip processing files when the (inexact) row count is above the fetch/limit:
I am pretty sure this is only valid if the estimate of the number of rows is conservative (aka the real value is at least as large as the statistics), but the code uses the value even for
Precision::Inexact:Another example is
ColumnStatistics::is_singleton, which I think is also only correct if the statistics are conservative (aka the actual min is no lower than the reported min and the actual max is no larger than the reported mx)Another challenge I found with using
Intervalas proposed in option 2 of #8078 (comment) is more mechanical but real:Intervalis in thedatafusion_physical_exprcrate butPrecisionis indatafusion_common, meaning I can't useIntervalinPrecisionwithout moving code around. Not impossible to do, but a data point.I will try prototyping option 1 sugested by @berkaysynnada #8078 (comment) and see what I find
I think there is a mismatch between the current use of
Precision::Inexactand its documentationSpecifically,
Precision::Inexactappears to be treated as a "conservative estimate" in several places (which is actually what we need in IOx), so perhaps another option would be to renamePrecision::ExacttoPrecision::Conservativeand document that it is a conservative estimate of the actual value 🤔I think the documentation is correct but you have found 2 bugs in the code:
-
Here, we must ensure the values are exact.
https://github.com/apache/arrow-datafusion/blob/87aeef5aa4d8f38a8328f8e51e530e6c9cd9afa9/datafusion/core/src/datasource/statistics.rs#L69-L73 -
Singleton check must be for exact values only.
https://github.com/apache/arrow-datafusion/blob/c3430d71179d68536008cd7272f4f57b7f50d4a2/datafusion/statistics/src/statistics.rs#L281-L286
I will open a PR now fixing these issues.
-
Another challenge I found with using
Intervalas proposed in option 2 of #8078 (comment) is more mechanical but real:Intervalis in thedatafusion_physical_exprcrate butPrecisionis indatafusion_common, meaning I can't useIntervalinPrecisionwithout moving code around. Not impossible to do, but a data point.I will try prototyping option 1 sugested by @berkaysynnada #8078 (comment) and see what I find
Just so you know, I have completed the interval refactor and submitted it for our internal review. After that, interval arithmetic will be in
datafusion_expr, and there will be no significant obstacles to moving it todatafusion_common.Reacted by Andrew LambAny thoughts on @berkaysynnada's proposal; i.e.
enum PointEstimate { Exact(T), Inexact(T), Absent, } enum Precision { PointEstimate, Range(PointEstimate, PointEstimate) }I think extending the
Precisionstruct like this, along with bugfixes like #8094, will address #8099.BTW this exercise revealed that the name
Precisionmaybe is not the best name. Maybe we should just useEstimate, which could either be aPointEstimateor aRange, which is tuple ofPointEstimates for bounds.Any thoughts on @berkaysynnada's proposal; i.e.
I really like it.
My plan was spend time trying to code up some of these proposals and see if it can fit in the existing framework (and thus potentially reveal in what cases Inexact would be used), etc
However, I spent much more of my time this week on apache/arrow-rs#5050 and related discussions around #7994 and #8017 than I had planned, so I am sadly behind
I hope to make progress tomorrow
Reacted by Berkay Şahin17 remaining items
let input_intervals: Vec<&Interval> = ....; // wrap input intervals with Statistics let temp_statistics = Statistics::new_from_bounds(&input_intervals); // compute output column statistics let output_column_statistics = expr.column_statistics(&temp_statistics)?; // use the output value if it was known let output_interval. = match output_column_statstics.value() { Precision::Absent | Precision::PointEstimation => None, Precision::Interval(interval) => interval };
This usage is weird in contexts where the concept of statistics isn't even applicable. I agree that it will work, but only so because we are forcing. I think the right pattern is to have
column_statisticssimply useevaluate_boundsas a subroutine for computing hard bounds (and any other information of type (2) in my comment). For example, it would be quite natural for us to have have things likeevaluate_probability(not a great name!) that also takes expressions and does some sort of probabilistic computation (maybe PDF related). Then,column_statisticswould also useevaluate_probabilityas a subroutine.So I see something like
column_statisticsas a more general API that defers to lower level APIs to collect information.Reacted by Andrew LambFWIW there is a current move to add statistics into the arrow format itself:
I actually think we could standardize on converting to/from that format.
I don't think the arrow proposal handles the subtelty about intervals, known facts, etc but we should at least be aware of them (thanks @edmondop for pointing this out to me)
I don't think the arrow proposal handles the subtelty about intervals, known facts, etc but we should at least be aware of them (thanks @edmondop for pointing this out to me)
Definitely. We will probably leverage that information to construct our own
Statistics.BTW we will start working on this tomorrow.
Reacted by Matthew Cramerus and Sasha SyrotenkoReacted by Andrew Lamb and Adrian Garcia BadaraccoI don't think the arrow proposal handles the subtelty about intervals, known facts, etc but we should at least be aware of them (thanks @edmondop for pointing this out to me)
Definitely. We will probably leverage that information to construct our own
Statistics.BTW we will start working on this tomorrow.
@ozankabak sorry to bump you, I'm looking into this stuff and see that your comment is from 2024 and there are a dozen linked issues / PRs. Is there a summary of progress that could be posted here?
No worries - we didn't have time to push it forward after merging the evaluation/propagation mechanism for statistics. I am not aware a summary document, I think the easiest way to get a hang of what is done is to look at git blame and identify the relevant PRs (there are not that many).
Reading through the issues and posting my thoughts as I go. I am particularly interested in improving the
Statisticsthat gets attached to files and partitions:It seems that just hasn't been updated to use
Distributioninstead ofPrecision. Doing this requires a re-design of theStatisticsstruct and handling all of the breaking changes. I think v50 already has a lot of breaking changes so we should not try to put it into this release, but maybe v51. I have some ideas for other changes as well (namely: instead of requiring aColumnStatisticsfor each column even those that are not present we can only include them for those that are somehow, otherwise a lot of memory is required for wide tables, it's fine forSchemabut this structure exists once per file).@alamb @ozankabak let me know if that sounds correct
Doing this requires a re-design of the Statistics struct and handling all of the breaking changes.
Yes I think that is the biggest challenge
What makes it hard is that https://docs.rs/datafusion/latest/datafusion/common/struct.Statistics.html
has a bunch of public fields
pub struct Statistics { pub num_rows: Precision<usize>, pub total_byte_size: Precision<usize>, pub column_statistics: Vec<ColumnStatistics>, }
I think we can probably prepare for / take most of the API pain at first by changing this so the fields are not pub
pub struct Statistics { num_rows: Precision<usize>, total_byte_size: Precision<usize>, column_statistics: Vec<ColumnStatistics>, }
And then adding appropriate APIs to construct / check statistics
Like
let statistics = Statistics::builder(( .with_num_rows(Precision::Exact(20)) .build();
I gave it a shot in #17980.
I decided to create a new struct to replaceColumnStatisticsand a new struct to replaceStatisticsand added them as fields toPartitionedFile, making it backwards compatible to at least experiment with the new thing.A couple interesting things I found:
- Have we considered
Precision<Distribution>? It seems to nicely encapsulate what I think we want. - I included and think we want several other public fields, e.g. I think
total_byte_size: Precision<usize>is sensible and we should keep it. The main thing to replace wasmin: Precision<ScalarValue>andmax: Precision<ScalarValue>. I keptsum: Precision<ScalarValue>but feel that maybe that could just be made to be a field of theDistributionvariants? @berkaysynnada any thoughts? - I couldn't put the new structs alongside the old ones because
datafusion/commondoes not depend ondatafusion/expr-common. Seems like we might need to do a bit of juggling if we want to accessDistributionfrom the current location. - I don't love the current reperesentation
Vec<ColumnStatistics>for several reasons: (1) for wide tables / scans having emptyColumnStatisticstakes up a lot of memory. We had a situation where we essentially initialized every file in the table (say 10k) and initialized empty stats for each column (say 100) so we ended up with a couple hunded MB of empty stats; (2) it depends on the projection so it's somewhat hard to tie statistics to the column by name; (3) when you apply a projection you have to depp clone all of the statistics. I tried to address (1) and (3) by changing the representation toVec<Option<Arc<ColumnDistributionStatistics>>which makes it cheaper to initailize as empty and cheaper to project. But that doesn't address (2). Should we also include the table schema for each file? It's just one moreArc'dfield. Then we could have helpers to get projected statistics for certain columns or certain indices, etc.
- Have we considered
I gave it a shot in #17980.
I think you mean
Have we considered Precision? It seems to nicely encapsulate what I think we want.
Seems like we might need to do a bit of juggling if we want to access Distribution from the current location.I can't help but feel the current
Distributiondoesn't have many practical benefits -- specifically the idea of having mathemetical descriptions of value distributions is intellectually appealing, but I have never see actual query engines use it (because real data is never completely described by those theoretical distributions). Maybe I am missing somethingI don't love the current reperesentation Vec for several reasons:
I broadly agree with this concern
But that doesn't address (2).
Maybe the statistics could include a
SchemaRefpointer so projections could be converted back into names when neededI tried to address (1) and (3) by changing the representation to Vec<Option<Arc> which makes it cheaper to initailize as empty and cheaper to project.
I think that is basically the right internal representation, but in terms of software engineering, I personally recommend:
Updating statistics initially to be a wrapper:
/// Not sure of this name. I am open to new ones /// StatisticsV2 was proposed at some point /// but that seems to offer no additional context. pub struct RelationStatistics { inner: Vec<Option<Arc<ColumnStatistics>> }
And then there should be an easy way to convert between
Vec<ColumnStatistics>to a RelationStatistics to ease migrationThen we could change the signatgure of ExecutionPlan to retun an arc'd version
trait ExecutionPlan { /// retun the statistics for this plan pub fn relation_statistics(&self) -> &Arc<RelationStatistics>; ... }
Then copying entire nodes would be fast (copy an Arc) as would projecting statistics 🤔
I can't help but feel the current Distribution doesn't have many practical benefits -- specifically the idea of having mathemetical descriptions of value distributions is intellectually appealing, but I have never see actual query engines use it (because real data is never completely described by those theoretical distributions). Maybe I am missing something
FWIW I do agree with this. For example take the range of values. Currently that's two different stat values
minandmax, but that should probably be encapsulated using 1 struct/enum. ForDistribution<ScalarValue>, it's unlikely that you ever gonna use anything else thanGenericbecause parquet -- or most other data sources -- give us it's really only a range with inclusive or exclusive bounds. So the entire enum is mostly unused.Then if we look at
GenericDistributionand it's constructor the issue is again that it requires knowledge like variance, median, and mean, which you likely never gonna know for most data sources. In fact if you have any filtered data source, then calculating themedianis virtually impossible if you wanna do anything that is remotely performant. So that's another 75% of the interface gone/unusable.So what's kinda left is the
Intervaltype and the kinda nice API methods around it. So maybe we could use that?I also feel that there's a slight conflict of interest or at least two camps here:
- statistics always-correct optimizers: Some people use statistics for optimizers like join ordering. There a wrong statistics often only results in slower execution, but never wrong results. That is kinda reflected in a lot of statistics calculation in the DF code base.
- correctness: Some plan transformers (InfluxData for example has one) rely on the statistics that actually can make hard promises, i.e. "all values are FOR SURE in this range". In that case, you really wanna be picky about what the stats do.
Reacted by xudong.w, Andrew Lamb and Alessandro SolimandoOne use case for
DistributionI wanted to explore that is compatible with Parquet is what I'll call a "footer table sample". I don't remember where I heard of this first or what I should call it, but I did discuss it with Hannes of DuckDB and it sounds like a really cool idea. TLDR is it's expensive to randomly sample compressed columnar storage like Parquet, but if you store a pre-sampled portion of the file e.g. as the last row group you can get very good estimates for all kinds of things (filter selectivity, cardinality of any column, etc.) and it's very efficient IO-wise to get that data (it's all nicely packed into 1 read unit). My thought is that something like this could be used to easily get estimated distributions and cardinality from the data.I also feel that there's a slight conflict of interest or at least two camps here:
- statistics always-correct optimizers: Some people use statistics for optimizers like join ordering. There a wrong statistics often only results in slower execution, but never wrong results. That is kinda reflected in a lot of statistics calculation in the DF code base.
- correctness: Some plan transformers (InfluxData for example has one) rely on the statistics that actually can make hard promises, i.e. "all values are FOR SURE in this range". In that case, you really wanna be picky about what the stats do.
I agree with this. My biggest issue with the current statistics is that we only have
ExactandInexactbutInexactisn't really what you want for the second case you list, you want something likeBounded.I also think the current statistics is lacking info like the size of each column which is much better than the total file size in almost every use case (most queries are not
select *).Reacted by Dan KingTLDR is it's expensive to randomly sample compressed columnar storage like Parquet, but if you store a pre-sampled portion of the file e.g. as the last row group you can get very good estimates for all kinds of things (filter selectivity, cardinality of any column, etc.)
Sounds like a good usecase for https://datafusion.apache.org/blog/2025/07/14/user-defined-parquet-indexes/ 🤔
TLDR is it's expensive to randomly sample compressed columnar storage like Parquet, but if you store a pre-sampled portion of the file e.g. as the last row group you can get very good estimates for all kinds of things (filter selectivity, cardinality of any column, etc.)
Sounds like a good usecase for https://datafusion.apache.org/blog/2025/07/14/user-defined-parquet-indexes/ 🤔
Yep could be that! I was thinking maybe the last row group would be beneficial because (assuming the data is basically Parquet data) it avoids having to re-encode all of the metadata. Also sadly our Parquet reader cannot be pointed at a byte range of a file (I think that'd be easy to fix in a PR). But it does make it incompatible with other readers (if they don't know to skip the last row group...). Anyway that should be discussed in another thread I just wanted to share the idea as a possibly use case for more advanced statistics.
Yep could be that! I was thinking maybe the last row group would be beneficial because (assuming the data is basically Parquet data)
This would work well if the data isn't sorted before writing (so the footer is a reasonably proxy for a random sample). If you sort the data beforehand the last row group probably isn't a good random sample
Also sadly our Parquet reader cannot be pointed at a byte range of a file (I think that'd be easy to fix in a PR)
With the metadata you can always figure out the ranges of each column chunk
However, I don't think you can just get the last 10% of the rows in the last row group, because data is stored column by column, so the data for the last 10% of the rows are going to be spread across multiple distinct ranges (for each column)
My thought was to take a random sample after sorting (thus duplicating random rows). Then make that its own row group. This means you can also use the table sample as a proxy for the sortedness of the data. Then you don't have to read a fraction of a row group. You'd read the metadata (which is then cached), make an access plan that reads only the last row group and you've got your table sample. Then for the actual scan you'd initialize the access plan with
Skip(row_group_count)so you always skip the sample row group. This could happen at the same time statistics are collected.
Is your feature request related to a problem or challenge?
This has come up a few times, most recently in discussions with @berkaysynnada on apache/arrow-rs#5037 (comment)
Usecase 1 is that for large binary/string columns, formats like parquet allow storing a truncated value that does not actually appear in the data. Given that values are stored in the min/max metadata, storing truncated values keeps the size of metadata down
For example, for a string column that has very long values, it requires much less space to store a short value slightly lower than the actual minimum as the "minimum" statistics value, and one that is slightly higher than the actual maximum as the "maximum" statistics value.
For example:
aaa......zqqq......qarThere is a similar usecase when applying a Filter, as described by @korowa on #5646 (comment) and we have a similar one in IOx where the operator may remove values, but won't decrease the minimum value or increase the maximum value in any column
Currently
Precisiononly representsExactandInexact, there is no way to represent "unexact, but bounded above/below"Describe the solution you'd like
Per @berkaysynnada I propose changing
Precision::Inexactto a new variantPrecision::Betweenwhich would store anIntervalof known min/maxes of the value.This is a quite general formulation, and it can describe "how" inexact the values are.
This would have the benefit of being very expressive (Intervals can represent open/closed bounds, etc)
Describe alternatives you've considered
There is also a possibility of introducing a simpler, but more limited version of these statistics, like:
Additional context
No response