apache / apache/datafusion

Control ordering of file opened based on statistics

Open
#17,271 3 comments 1 reaction 1 assignee Claimed by @adriangb View on GitHub
Dominant language
Rust
Stars
9.3k
Forks
2.4k
Avg merge
3d 7h
Merged PRs (30d)
344

Description

In particular for a query that has `ORDER BY col ASC LIMIT 10` we will use a TopK operator + dynamic filter pushdown to prune files. This can over >20x faster query performance but effectiveness depends largely on the order in which files are opened: in the pathological case that the files with the largest `col` are opened first for we'll have to open every single file. The ideal case would be that we open files with the smallest `col` first and in the first file we find the 10 smallest `col` and thus are able to skip all others based on statistics. This also causes less churn in the TopK heap, etc.

Currently ListingTable orders files based on a known sort order if provided (see https://github.com/apache/datafusion/issues/4177) or their path: https://github.com/apache/datafusion/blob/3b7eb267ddbac10136b395989a76acfe9664425e/datafusion/datasource/src/file_groups.rs#L444-L447

I'd like to propose that instead we pass down the _preferred_ sort order for the query (instead of a hardcoded known sort order) and try to use statistics to sort the files _within each partition/group_ to best match that sort order. I think that covers both the Influx IOx use cases and more general use cases where strict non-overlapping ordering is not required but generally ordering the file opens to agree with sort operators is beneficial.

I believe the main barrier to this is that sort information is not passed down into `TableProvider`. It would have to be an additional option to `scan`. Adding an additional option would be a breaking change that impacts a lot of users and `scan` already has a lot of option. Hence I propose the following:

```rust
struct ScanOptions {
preferred_ordering: Vec,
filters: Vec,
limit: Option,
}

struct ScanResult {
/// The ExecutionPlan to run.
plan: Arc,
// Remaining filters that were not completely evaluated during `scan_with_options()`.
filters: Vec,
}

trait TableProvider {
fn scan_with_options(&self, options: ScanOptions) -> Result;
#[deprecated]
fn scan(&self, ...) -> Result>;
#[deprecated]
fn supports_filters_pushdown(&self, filters: &[&Expr]) -> Result> { ...
}
```

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.