apache / apache/datafusion

Improve performance of large sorts with Cascaded merge / tree

Open
#7,181 1 comment 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?

While working on https://github.com/apache/arrow-datafusion/pull/7179 I noticed a potential improvement

The key observation is that merging `K` sorted streams of total rows `N`:
1. takes time proportional to `O(N*K)`
2. is a single threaded operation

`K` is often called the "Fan In" of the merge

The implementation of [ExternalSorter::in_mem_sort_stream](https://github.com/apache/arrow-datafusion/blob/1e9816b07edd8e73bb45ababa20a7c22de492d96/datafusion/core/src/physical_plan/sorts/sort.rs#L259-L297) will effectively merge all the
buffered batches at once, as shown below. This can be a very large fan in -- 100s or 1000s of RecordBatches

```text
┌─────┐ ┌─────┐
│ 2 │ │ 1 │
│ 3 │ │ 2 │
│ 1 │─ ─▶ sort ─ ─▶│ 2 │─ ─ ─ ─ ─ ─ ┐
│ 4 │ │ 3 │
│ 2 │ │ 4 │ │
└─────┘ └─────┘
┌─────┐ ┌─────┐ │
│ 1 │ │ 1 │
│ 4 │─ ▶ sort ─ ─ ▶│ 1 ├ ─ ┐ │
│ 1 │ │ 4 │
└─────┘ └─────┘ │ │
... ... ▼

Could be 100s depending on ─ ─▶ merge ─ ─ ─ ─ ─▶ sorted output
the data being sorted stream

... ...
┌─────┐ ┌─────┐ │
│ 3 │ │ 3 │
│ 1 │─ ▶ sort ─ ─ ▶│ 1 │─ ─ ─ ─ ─ ─ ┤
└─────┘ └─────┘
┌─────┐ ┌─────┐ │
│ 4 │ │ 3 │
│ 3 │─ ▶ sort ─ ─ ▶│ 4 │─ ─ ─ ─ ─ ─ ┘
└─────┘ └─────┘

in_mem_batches

```

### Describe the solution you'd like

A classical approach to such sorts is to use a "cascaded merge" which uses a series of merge operations each with a limited the fanout (e.g. to 10)

```
┌─────┐ ┌─────┐
│ 2 │ │ 1 │
│ 3 │ │ 2 │
│ 1 │─ ─▶ sort ─ ─▶│ 2 │─ ─ ─ ─ ─ ─ ─ ─ ┐
│ 4 │ │ 3 │
│ 2 │ │ 4 │ │
└─────┘ └─────┘
┌─────┐ ┌─────┐ ▼
│ 1 │ │ 1 │
│ 4 │─ ▶ sort ─ ─ ▶│ 1 ├ ─ ─ ─ ─ ─ ▶ merge ─ ─ ─ ─
│ 1 │ │ 4 │ │
└─────┘ └─────┘
... ... ... ▼

merge ─ ─ ─ ─ ─ ─ ▶ sorted output
stream

... ... ... │
┌─────┐ ┌─────┐
│ 3 │ │ 3 │ │
│ 1 │─ ▶ sort ─ ─ ▶│ 1 │─ ─ ─ ─ ─ ─▶ merge ─ ─ ─ ─
└─────┘ └─────┘
┌─────┐ ┌─────┐ ▲
│ 4 │ │ 3 │
│ 3 │─ ▶ sort ─ ─ ▶│ 4 │─ ─ ─ ─ ─ ─ ─ ─ ┘
└─────┘ └─────┘

in_mem_batches do a series of merges that
each has a limited fan-in
(number of inputs)
```

This is often better because:
1. Is `O(N*ln(N)*ln(K))` , there is some additional overhead of `ln(N)` as the same row must now be compared several times
2. the intermediate merges can be run in parallel on multiple cores (though the final one is still single threaded)

It would be awesome if someone wanted to:

1. Verify the theory that there is a large fan in for large sorts
3. Implement a cascaded merge and measure if it improves performance

The sort [benchmark](https://github.com/apache/arrow-datafusion/tree/main/benchmarks) (TODO) (thanks @jaylmiller!) may be interesting:

```
cargo run --release --bin parquet -- sort --path ./data --scale-factor 1.0
```

### Describe alternatives you've considered

Another potential variation might be to get more cores involves in the merging by parallelizing the merge, as described in the [Morsel-Driven Parallelism paper](https://db.in.tum.de/~leis/papers/morsels.pdf)

Screenshot 2023-08-02 at 8 56 46 AM

### Additional context

_No response_

Contributor guide

Open the contributing guide

Research direction

Start with ExternalSorter::in_mem_sort_stream in datafusion/core/src/physical_plan/sorts/sort.rs and verify how many buffered batches are merged for large sorts. Run the parquet benchmark with `cargo run --release --bin parquet -- sort --path ./data --scale-factor 1.0`, then implement and measure a cascaded merge; done means the fan-in and performance impact are documented by benchmark results.

Written by the indexing model from the issue text.

Assessment

Tech stack
rust
Domain
data-engineering, performance
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.