Skip to content

benchmark partition-aware dynamic filters - #3

Closed
jayshrivastava wants to merge 0 commit into
js/benchmark-partition-aware-dynamic-filters-base-shafrom
js/benchmark-partition-aware-dynamic-filters
Closed

benchmark partition-aware dynamic filters#3
jayshrivastava wants to merge 0 commit into
js/benchmark-partition-aware-dynamic-filters-base-shafrom
js/benchmark-partition-aware-dynamic-filters

Conversation

@jayshrivastava

@jayshrivastava jayshrivastava commented Jul 16, 2026

Copy link
Copy Markdown
Owner

🚨 Disclaimer

  • The code was entirely generated by codex and was not reviewed
  • The benchmarking methodology was human-revised
  • The PR description is human-revised

Background

For HashJoinExec mode=Partitioned, we use CASE expressions to achieve partition-aware dynamic filtering.

CASE hash(expr) % num_partitions
  WHEN 0 THEN filter_expr_for_partition_0
  WHEN 1 THEN filter_expr_for_partition_1
  WHEN 2 THEN filter_expr_for_partition_2
   ...
  ELSE false
END

This case feature was added in apache#18451, but theres' no benchmarks.

Does this necessarily perform better than something like this?

filter_expr_for_partition_0
  OR filter_expr_for_partition_1
  OR filter_expr_for_partition_2
  ...

Goal

Benchmark how significantly this CASE expression actually improves query performance.

Setup

Benchmark and Data

  • Benchmark: dfbench hj (query 23 and new "query 24"

  • Dataset: TPC-H SF1 parquet generated with tpchgen-cli v3.0.0

  • Profile: cargo run --profile release-nonlto

  • Iterations: 5

  • Host: Linux aarch64

  • Query 23 uses a partitioned hash join, so we selected it

    • The probe-side filter is over a synthetic string expression so row group pruning is not very effective. For this reason, we added query 24
  • Query 24 is a custom query which uses a stored column as the join key, allowing more effective pruning

Benchmark Cases

For all cases, force partitioned joins

SET datafusion.optimizer.hash_join_single_partition_threshold = 0;
SET datafusion.optimizer.hash_join_single_partition_threshold_rows = 0;
  • We mainly want to compare B vs C and D vs E (CASEexpression vs global OR expression).
  • B&C turn pruning off because it is does not support CASE expressions, making the comparison a bit more fair. When pruning is not supported, we read all the row groups and apply dynamic filters on the rows after reading.
Case Session properties
A: baseline no dyn filter - datafusion.optimizer.enable_dynamic_filter_pushdown=false
- datafusion.execution.parquet.pruning=false
- datafusion.execution.parquet.enable_page_index=false
- datafusion.execution.parquet.bloom_filter_on_read=false
- datafusion.execution.parquet.pushdown_filters=true
B: case expr no pruning - datafusion.optimizer.enable_dynamic_filter_pushdown=true
- datafusion.optimizer.hash_join_dynamic_filter_partitioned_expr_style=case
- datafusion.execution.parquet.pruning=false
- datafusion.execution.parquet.enable_page_index=false
- datafusion.execution.parquet.bloom_filter_on_read=false
- datafusion.execution.parquet.pushdown_filters=true
C: global_or no pruning - datafusion.optimizer.enable_dynamic_filter_pushdown=true
- datafusion.optimizer.hash_join_dynamic_filter_partitioned_expr_style=global_or
- datafusion.execution.parquet.pruning=false
- datafusion.execution.parquet.enable_page_index=false
- datafusion.execution.parquet.bloom_filter_on_read=false
- datafusion.execution.parquet.pushdown_filters=true
D: case full - datafusion.optimizer.enable_dynamic_filter_pushdown=true
- datafusion.optimizer.hash_join_dynamic_filter_partitioned_expr_style=case
- datafusion.execution.parquet.pruning=true
- datafusion.execution.parquet.enable_page_index=true
- datafusion.execution.parquet.bloom_filter_on_read=true
- datafusion.execution.parquet.pushdown_filters=true
E: global_or full - datafusion.optimizer.enable_dynamic_filter_pushdown=true
- datafusion.optimizer.hash_join_dynamic_filter_partitioned_expr_style=global_or
- datafusion.execution.parquet.pruning=true
- datafusion.execution.parquet.enable_page_index=true
- datafusion.execution.parquet.bloom_filter_on_read=true
- datafusion.execution.parquet.pushdown_filters=true

Results

The case expression has no benefit.

Q23 Default Partitions

Case Avg (ms)
A: baseline no dyn filter 22.983
B: case expr no pruning 25.309
C: global_or no pruning 24.467
D: case full 24.886
E: global_or full 27.777

Q23 Four Partitions

Case Avg (ms)
A: baseline no dyn filter 34.524
B: case expr no pruning 36.256
C: global_or no pruning 35.125
D: case full 35.267
E: global_or full 35.206

Q24 Default Partitions

Case Avg (ms)
A: baseline no dyn filter 26.006
B: case expr no pruning 60.226
C: global_or no pruning 31.843
D: case full 49.047
E: global_or full 21.988

Q24 Four Partitions

Case Avg (ms)
A: baseline no dyn filter 48.665
B: case expr no pruning 81.794
C: global_or no pruning 35.737
D: case full 76.957
E: global_or full 7.874
  • Q23 doesn't benefit from dynamic filtering, so let's ignore that
  • Q24
    • global or expression is strictly better than the case expr with and without row-level pruning

Why is CASE so slow?

  • (pruning on) Q24 partitions=4 D vs E, pushdown_rows_pruned=19.03 K for the global OR vs only pushdown_rows_pruned=6.00 M for CASE. Avoiding CASE allows us to prune row groups.
  • (pruning off) Q24 partitions=4 B vs C pushdown_eval_time=21.16ms for the global OR vs pushdown_eval_time=196.13ms for CASE. The CASE expr uses more compute.

@jayshrivastava
jayshrivastava changed the base branch from main to js/benchmark-partition-aware-dynamic-filters-base-sha July 16, 2026 17:21
@jayshrivastava jayshrivastava changed the title Js/benchmark partition aware dynamic filters benchmark partition-aware dynamic filters Jul 16, 2026
@jayshrivastava
jayshrivastava force-pushed the js/benchmark-partition-aware-dynamic-filters branch from 80d2214 to c1647e3 Compare July 17, 2026 16:33
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant