Out-Of-Core Shuffling w/ RapidsMPF

Shuffling data at 1.8 TiB/s! RapidsMPF is a reusable, out-of-core shuffler that turns shuffling OOM headaches into a spill you can budget for.

Shuffling is the crux of structured data analytics distributed or otherwise. It's a core component of key data operations like: join, groupby, merge, sort, etc. A full distributed shuffle can move all the data from every process to every other process, an all-to-all. This is very costly and many sophisticated techniques have been developed to avoid this operation as much as possible. Shuffles aren't particularly computationally challenging: calculating the hashes to route data is fairly cheap. They are nonetheless expensive in a workflow, for a variety of reasons: