apache / apache/datafusion

Optimizations for hash join operator

Open
#18,942 1 comment 0 reactions 0 assignees View on GitHub
performance
Dominant language
Rust
Stars
9.3k
Forks
2.4k
Avg merge
3d 7h
Merged PRs (30d)
344

Description

> Thanks! Besides looking at optimizing the join order during planning time or dynamic (I think there are a couple of issues covering that), we can look at what makes the operator slow in more challenging scenario's.
>
> Some optimizations for the current operator come to mind that might improve the current hash join operator in certain scenario's, while keeping the same algorithm:
>
> * Reuse the allocation of `Vec` indices between calls. This probably helps when the amount of matching indices is low (compared to the batch size).
> * (Related): Keep building matching indices until `limit` rows have been reached and use `interleave` to collect the batches. That probably makes the operator more cache efficient as accessing the map / chain is done at the same time, before producing output batches from the input data. This also helps with avoiding the overhead of a later `CoalesceBatches`, which might help as well.
> * Instead of building indices for the right side, we can build a boolean mask / filter to mark match / no match. This reduces memory usage (somewhat) plus a boolean filter is much faster for low selectivity (i.e. most of the right side matches). We then should use the coalesce kernel to produce the right side arrays.
>
> I opened https://github.com/apache/datafusion/issues/18939 for exploring to use a different algorithm (radix hash joins), which additionally should improve the performance of our join operators by making the algorithm more cache efficient.

_Originally posted by @Dandandan in [#17494](https://github.com/apache/datafusion/issues/17494#issuecomment-3581011421)_

Contributor guide

Open the contributing guide

Research direction

The issue names no files or tests. Begin by locating the current hash join operator, then review the three proposed optimization directions and related issue #18939. Done would require selecting an approach and showing improved performance in the relevant scenarios without changing the algorithm.

Written by the indexing model from the issue text.

Assessment

Tech stack
rust
Domain
databases
Issue type
Refactor
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Needs clarification
Newbie friendliness
25/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.