hn.today

Accelerated Out of Core Shuffling

quasiben.github.io11 points0 comments
Screenshot of Accelerated Out of Core Shuffling

Shuffling is presented as the critical bottleneck in distributed structured-data operations like joins, groupbys and sorts because it forces all-to-all data movement and often requires materializing many copies of tables in memory. A classic in-memory hash join has a build phase and a probe phase, but distributed joins must first route rows by hashing join keys so that matching keys land on the same rank. That routing, packing, sending and waiting for all producers creates large memory pressure and synchronization barriers; minimizing full materialization and enabling streaming, spillable shuffles is therefore essential. RapidsMPF was developed to address this: a reusable shuffle library that supports out-of-core spilling and accelerated transport plus an actor network for streaming pipelines. The shuffle component is usable from C++ or Python and is already adopted by projects such as cuDF and Polars.

Benchmarks use bench_shuffle on a single DGXB200 (8 Blackwell GPUs with 180 GB VRAM each) via UCX/UCXX to measure shuffle performance and memory behavior. A representative run used 10 columns and 536,870,912 rows per rank (2 GiB per column → 20 GiB per rank, 160 GiB total across 8 ranks) and reported per-rank local throughput around 229-236 GiB/s and an aggregate reported global throughput near 1.8 TiB/s (noting a reporting bug that multiplies local throughput by rank count). Memory profiling shows device memory peaks around 60 GiB with alloc-device ~17.5 GiB and shuffle payload send/recv ~17.5 GiB, demonstrating high throughput while supporting out-of-core operation and controlled spilling.

Read on quasiben.github.io0 comments on Hacker News

Summary generated by AI from the linked article. hn.today is not affiliated with Hacker News or Y Combinator.

More in Security

The daily digest

Today's best Hacker News stories, summarized and screenshotted, one email a day.