apache / apache/datafusion

Improve performance of large sorts with Cascaded merge / tree

Aperta
#7,181 1 commento 0 reazioni 0 assegnatari Vedi su GitHub
enhancement
Lingua principale
Rust
Stelle
9.3k
Fork
2.4k
Merge medio
3g 11h
PR unite (30g)
360

Descrizione

### 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_

Guida per i contributori

Apri la guida per i contributori

Direzione di ricerca

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.

Scritto dal modello di indicizzazione a partire dal testo della issue.

Valutazione

Stack tecnologico
rust
Ambito
data-engineering, performance
Tipo di issue
Funzionalità
Difficoltà
5/5
Tempo stimato
Più di una settimana
Stato di attività
Ferma
Chiarezza
Abbastanza chiara
Idoneità per principianti
35/100

Ricevi le nuove issue nella tua casella

Un breve riepilogo di issue GitHub adatte ai principianti.