Out-Of-Core Shuffling w/ RapidsMPF
Summary
This article presents Out-Of-Core Shuffling w/ RapidsMPF, focusing on memory-efficient, GPU-accelerated data shuffles used in joins and ETL workflows. It explains why shuffling is memory-intensive, describes the RapidsMPF architecture (shuffle library and streaming actor network), and surveys benchmarking results on a DGX B200 that reach ~1.8 TiB/s global throughput. It also discusses memory pressure, spilling strategies, and host-device transfers (including the impact of PCIe bandwidth and UCXX) and finishes with practical takeaways for designing scalable, out-of-core data pipelines.