Skip to content
HN On Hacker News ↗

Out-Of-Core Shuffling w/ RapidsMPF - Benjamin Zaitlen

▲ 11 points • 0 comments • by quasiben • 2w ago • HN discussion ↗

Pangram verdict · v3.3

We believe that this entire text is human-written.

5 %

AI likelihood · overall

Human
100% human-written 0% AI-generated
SEGMENTS · HUMAN 1 of 1
SEGMENTS · AI 0 of 1
WORD COUNT 1,561
PEAK AI % 5% · §1
Analyzed
Sep 27
backend: pangram/v3.3
Segments scanned
1 windows
avg 1561 words each
Distribution
100 / 0%
human / AI fraction
Verdict
Human
Pangram v3.3

Article text · 1,561 words · 1 segments analyzed

Human AI-generated
§1 Human · 5%

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: Memory intensive: shuffles can require holding onto a full copy of all the data or, in streaming cases, memory pressure can build and cause OOMs. Transport: The data physically has to be moved from Process A->Process B or Node A->Node B so it can only move at speed of the transport layer. Synchronization: Output data cannot be consumed until all producers have finished contributing data. In a bulk-synchronous engine, this barrier can stall the entire execution plan. Because shuffling is hard, slow, memory-intensive, and critical, it's historically where RapidsMPF started. Why do Joins need Shuffles?¶ A quick primer on joining tables. If we have two tables: partsupp and lineitem and we want to join them, what happens? partsupp.join( lineitem, left_on=["ps_partkey", "ps_suppkey"], right_on=["l_partkey", "l_suppkey"], ) # or SELECT * FROM partsupp JOIN lineitem ON partsupp.ps_partkey = lineitem.l_partkey AND partsupp.ps_suppkey = lineitem.l_suppkey In-Memory Joins¶ Inner joins are composed of two phases: build phase: the smaller table (partsupp) is scanned and a hash table is constructed over the join keys (ps_partkey, ps_suppkey), mapping each hashed key to the row it came from. This hash table is fully populated before the probe phase can begin. probe phase: the larger table (lineitem) is scanned and each row's keys (l_partkey, l_suppkey) are hashed. Matches between the hash of the build and probe table join keys emit an output row combining columns from both tables. Rows without matches are dropped (technically, there's also hash collision handling here, but ignore that for now). note: left, right, and full outer joins use the same build/probe strategy with different rules for unmatched rows At minimum, this in-memory join holds three tables: the build table, the probe table, and the output table and the hash table built over the build side. Distributed Join¶ In the in-memory case, all the data is already colocated within the same memory space. That's no longer true once the tables are spread across many processes/nodes/ranks or tables are batched for "streaming" joins. A rank/process can only join rows in resident memory. Eventually, an in-memory join will occur, but first we'll need to get all the matching keys for the build and probe tables on the samerank. The cartoon graphic below represents how various rows of the same color are shuffled into the same output partition, and those partitions live on different ranks. To execute a distributed hash-join one must do the following: Scan the build table, hash the join keys of each row to pick a destination partition, hash(keys) % n_out_partitions. Pack and send each row to the rank that will own it. Scan the probe table and route the rows the same way so that probe keys land on the rank already holding the build keys with the same hash. Wait until every rank has finished sending. Only then is a rank guaranteed to hold every row, from both tables, for the keys it owns. Run the in-memory join from above on each rank's local slice: build a hash table over its build rows, probe it with its probe rows, emit matches. In the worst case, if every stage is fully materialized before the next one begins, a single rank is holding all of the following at once: build table (source scan) probe table (source scan) staged build table (packed for send) staged probe table (packed for send) shuffled build slice (received) shuffled probe slice (received) hash table over the build slice output table This is why shuffling is memory intensive rather than compute intensive. The hashing itself is cheap and it's why an out-of-core shuffle implementation should be a primary focus when setting out to build an ETL engine. Additionally, being able to stream data rather than fully materialize the tables before shuffling is critical for reducing memory pressure. For these reasons, we started RapidsMPF with the original goal of building a streaming out-of-core shuffler. RapidsMPF¶ RapidsMPF has expanded since its original conception. It is now a library composed of two large pieces: 1. A shuffle library designed for spilling / out-of-core memory handling, with accelerated transport 1. An actor network for constructing streaming data pipelines Users today can still adopt just the shuffling component of RapidsMPF (C++ or Python interfaces). We've seen this adoption in NeMo-Curator and, experimentally, in Ray Data. Most importantly, cuDF Polars uses RapidsMPF for both shuffles and the actor network. In a follow-up post we can dive into the actor network or if you're curious now I'd recommend reading the section on the streaming engine. Our shuffle implementation needs to: Be fast Scale Work with larger than VRAM (GPU) data (out-of-core) Be reusable and in the rest of this blog we'll focus our attention on shuffling under memory pressure. Benchmarking Setup¶ cuDF/RapidsMPF has an easy-to-use C++ benchmark: bench_shuffle which helps us study how the RapidsMPF shuffling implementation works across varied hardware: transports, number of GPUs, etc, as well as varied configuration like: input/output partition sizes, memory resources, etc. Here's a full breakdown of what the current bench_shuffle test exposes to users, along with the values I use throughout this post. Generally speaking, this benchmark builds tunable amounts of random 32-bit (4 byte) integers per rank (per GPU), shuffles all the data, and completes (there is no join here, just the shuffle). Flag Meaning Values used here -C <name> Communicator ucxx -c <n> Number of columns 10 -r <n> Number of timed runs 10 -w <n> Number of warmup runs 3 -n <n> Number of rows per rank 536870912 (2 GiB per column at 4 bytes/row) -p <n> Number of input partitions per rank 1 -o <n> Number of output partitions per rank 8 (one per rank) -m <name> RMM memory resource pool -l <n> Device memory limit in MiB omitted = unlimited (binary default -1); 32768 down to 12288 in the spill sweep -s Enable output discard (simulate streaming) flag, always set -x Enable memory profiling flag, always set -g Use pre-partitioned input tables flag, always set For all tests we are going to use a single DGXB200 and we are going to use rrun, an mpi like multiprocess launch tool capable of binding processes to NUMA nodes, to launch the shuffles. Simple Shuffling¶ A DGXB200 has 8 Blackwell GPUs, each with 180GBs of VRAM, and 2 Intel® Xeon® Platinum 8570 Processors. To get a baseline, we'll start by shuffling data which comfortably fits across all GPUs. rrun -n 8 --bind-to cpu --bind-to memory -x UCX_MAX_RNDV_RAILS=1 -x UCX_PROTO_ENABLE=y -x UCX_WARN_UNUSED_ENV_VARS=n libcudf_streaming_bench_shuffle -C ucxx -w 3 -r 10 -m pool -g -s -x -p 1 -o 8 -c 10 -n 536870912 Here we are warming up the benchmark 3 times, then running the benchmark 10 times. There are 536_870_912 rows (-n), 10 columns (-c), 1 input partition per rank (-p), and the data will be shuffled into 8 output partitions (-o). We are also using UCXX/UCX to enable accelerated transport/GPUDirect RDMA. 536_870_912 rows * 4 bytes (32-bit ints) = 2 GiB per column 10 columns * 2 GiB = 20 GiB / rank 8 ranks * 20 GiB = 160 GiB total # example output for 20GiB/rank [6:PRINT:0:2026-09-16 02:09:53.934930934] elapsed: 17.91 s | local throughput: 1.12 GiB/s | global throughput: 8.93 GiB/s (warmup run) [5:PRINT:0:2026-09-16 02:09:53.935046244] elapsed: 17.91 s | local throughput: 1.12 GiB/s | global throughput: 8.93 GiB/s (warmup run) [4:PRINT:0:2026-09-16 02:09:54.072443638] elapsed: 94.58 ms | local throughput: 211.47 GiB/s | global throughput: 1.65 TiB/s (warmup run) [2:PRINT:0:2026-09-16 02:09:54.072450366] elapsed: 89.15 ms | local throughput: 224.35 GiB/s | global throughput: 1.75 TiB/s (warmup run) [0:PRINT:0:2026-09-16 02:09:54.328272061] elapsed: 86.89 ms | local throughput: 230.19 GiB/s | global throughput: 1.80 TiB/s [1:PRINT:0:2026-09-16 02:09:54.328286361] elapsed: 85.38 ms | local throughput: 234.25 GiB/s | global throughput: 1.83 TiB/s [7:PRINT:0:2026-09-16 02:09:54.328417037] elapsed: 84.58 ms | local throughput: 236.47 GiB/s | global throughput: 1.85 TiB/s [4:PRINT:0:2026-09-16 02:09:54.328429082] elapsed: 85.15 ms | local throughput: 234.87 GiB/s | global throughput: 1.83 TiB/s [2:PRINT:0:2026-09-16 02:09:54.328554677] elapsed: 86.28 ms | local throughput: 231.80 GiB/s | global throughput: 1.81 TiB/s Each rank posts how much time it spent shuffling, and the local and global throughput. Already we can observe that warming up has some cost as it runs slower than the "official" run. At the end, the program returns the average values per rank of local/global throughput and summary statistics per rank for where time was spent: time in shuffle, time allocating memory, spilling (if any), etc. Note The global throughput varies from rank to rank, which is technically wrong and is a reporting bug. Rather than summing the local throughputs across ranks, the benchmark reports the global throughput as each rank's own local throughput multiplied by the number of ranks. On a DGXB200 it's 8 x local throughput. But it will suffice for now while the bug is resolved.