apache / apache/datafusion

Implement dynamic discrete pruning through a join

Open
#25,291 0 comments 0 reactions 0 assignees View on GitHub
enhancement
Dominant language
Rust
Stars
9.3k
Forks
2.4k
Avg merge
3d 7h
Merged PRs (30d)
344

Description

### Is your feature request related to a problem or challenge?

Whilst we currently do have bounds (min/max pruning) pushed down dynamically through a hash join, we don't have the more advanced discrete version of it wired up, even though all the primitives are in place.

Consequently, we're leaving some performance benefits on the table, especially when a qualifier on the build side filters some non-contiguous narrow sub-ranges in a broader contiguous range (e.g. every March in a 5-year interval).

### Describe the solution you'd like

Anything that would bring down the scanned row count in the following toy example in line with the build side:
```sql
> select version();
+--------------------------------------------+
| version() |
+--------------------------------------------+
| Apache DataFusion 55.0.0, aarch64 on macos |
+--------------------------------------------+
1 row(s) fetched.
Elapsed 0.002 seconds.

> copy (select i as k, random() as v from generate_series(0, 1999999) t(i))
to '/tmp/fact.parquet'
stored as parquet options ('format.max_row_group_size' '1000');
+---------+
| count |
+---------+
| 2000000 |
+---------+
1 row(s) fetched.
Elapsed 0.093 seconds.

> create external table fact stored as parquet location '/tmp/fact.parquet';
0 row(s) fetched.
Elapsed 0.010 seconds.

> create table dim as
select i as k from generate_series(0, 1999999) t(i) where i % 10000 < 200;
0 row(s) fetched.
Elapsed 0.014 seconds.

> explain analyze select count(*), sum(v) from fact join dim on fact.k = dim.k;
...
| | DataSourceExec: file_groups={12 groups: [[tmp/fact.parquet:0..1929817], [tmp/fact.parquet:1929817..3859634], [tmp/fact.parquet:3859634..5789451], [tmp/fact.parquet:5789451..7719268], [tmp/fact.parquet:7719268..9649085], ...]}, projection=[k, v], output_ordering=[k@0 ASC NULLS LAST], file_type=parquet, predicate=DynamicFilter [ k@0 >= 0 AND k@0 <= 1990199 AND hash_lookup ], dynamic_rg_pruning=eligible, pruning_predicate=k_null_count@1 != row_count@2 AND k_max@0 >= 0 AND k_null_count@1 != row_count@2 AND k_min@3 <= 1990199, required_guarantees=[], metrics=[output_rows=1.99 M, elapsed_compute=584.15µs, output_bytes=248.9 MB, output_batches=1.99 K, files_ranges_pruned_statistics=12 total → 12 matched, row_groups_pruned_statistics=2.00 K total → 1.99 K matched, row_groups_pruned_bloom_filter=1.99 K total → 1.99 K matched, page_index_pages_pruned=1.99 K total → 1.99 K matched, page_index_rows_pruned=1.99 M total → 1.99 M matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, bytes_scanned=22.34 M, file_open_errors=0, file_scan_errors=0, files_opened=12, files_processed=12, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, row_groups_pruned_dynamic_filter=0, predicate_cache_inner_records=0, predicate_cache_records=0, bloom_filter_eval_time=23.02ms, metadata_load_time=1.30ms, page_index_eval_time=7.02ms, row_pushdown_eval_time=36ns, statistics_eval_time=842.73µs, time_elapsed_opening=28.55ms, time_elapsed_processing=120.79ms, time_elapsed_scanning_total=1.19s, time_elapsed_scanning_until_data=9.18ms, output_rows_skew=1.64%, scan_efficiency_ratio=96.46% (22.34 M/23.16 M)] |
...
```

In particular, note that all 2 million rows are scanned, even though the `dim` table (build side) has only 40K values (co-located in 10% of fact's row groups) to join on.

### Describe alternatives you've considered

One might be tempted to play around with `datafusion.execution.parquet.max_in_list_size`/`datafusion.optimizer.hash_join_inlist_pushdown_max_distinct_values` and thus (ab)use `InListExpr` to achieve this, which seem to guard this path, but that is probably less than ideal for a number of reasons (with pushdown_filters=true this might prove to be counter-productive, EXPLAINs are polluted, etc.).

### Additional context

Related to https://github.com/apache/datafusion/issues/7955.

Contributor guide

Open the contributing guide

Research direction

Start by reproducing the provided EXPLAIN ANALYZE example and compare the existing dynamic min/max pruning through a hash join with the discrete-pruning primitives mentioned in the issue. Trace the hash-join dynamic filter and Parquet row-group pruning paths; done means the example avoids scanning irrelevant row groups and preserves correct join results without relying on InListExpr settings.

Written by the indexing model from the issue text.

Assessment

Tech stack
rust
Domain
databases, performance
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Active
Clarity
Mostly clear
Newbie friendliness
45/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.