AI・機械学習
超高速アウトオブコアシャッフリング
Accelerated Out of Core Shuffling (quasiben.github.io)
要約
この記事は、RapidsMPFライブラリを用いた、1.8 TiB/sという驚異的な速度でのアウトオブコア(メモリ外)シャッフリングについて解説しています。シャッフリングはデータ分析における結合、集計、ソートなどの基幹操作に不可欠ですが、メモリ使用量やデータ転送のオーバーヘッドが課題でした。RapidsMPFは、この課題を解決し、GPUメモリを超える大規模データセットでも効率的なシャッフリングを実現します。
全文翻訳
超高速アウトオブコアシャッフリング
1.8 TiB/sでのデータシャッフリング!RapidsMPFは、シャッフリングにおけるメモリ不足の悩みを、予算内で処理できる「スピル」に変える、再利用可能なアウトオブコアシャッフラーです。シャッフリングは、構造化データ分析の中核であり、分散処理か否かにかかわらず、結合、集計、マージ、ソートなどの主要なデータ操作のコアコンポーネントです。完全な分散シャッフルは、すべてのプロセスから他のすべてのプロセスへ全データを移動させる可能性があり、これは非常にコストがかかるため、この操作を可能な限り回避するための洗練された技術が多く開発されてきました。シャッフリング自体は計算的にそれほど難しくありません。データをルーティングするためのハッシュを計算するのは比較的容易です。しかし、様々な理由からワークフローにおいては高コストになります。
メモリ集約的:シャッフリングは、全データの完全なコピーを保持する必要がある場合や、ストリーミングケースではメモリ圧力が蓄積してメモリ不足(OOM)を引き起こす可能性があります。
輸送:データは物理的にプロセスAからプロセスB、またはノードAからノードBへ移動する必要があるため、輸送層の速度でしか移動できません。
同期:出力データは、すべてのプロデューサーがデータの提供を完了するまで消費できません。バルク同期エンジンでは、このバリアが実行プラン全体を停滞させる可能性があります。
シャッフリングは困難で、遅く、メモリ集約的であり、かつクリティカルであるため、RapidsMPFが最初に焦点を当てたのはこの部分でした。
結合はなぜシャッフルを必要とするのか?
テーブルの結合に関する簡単な説明です。2つのテーブル、partsuppとlineitemがあり、それらを結合したい場合、何が起こるでしょうか?
partsupp.join(
lineitem,
left_on=["ps_partkey", "ps_suppkey"],
right_on=["l_partkey", "l_suppkey"],
)
# または SELECT * FROM partsupp JOIN lineitem ON partsupp.ps_partkey = lineitem.l_partkey AND partsupp.ps_suppkey = lineitem.l_suppkey
インメモリ結合
内部結合は2つのフェーズで構成されます。
ビルドフェーズ:小さい方のテーブル(partsupp)をスキャンし、結合キー(ps_partkey, ps_suppkey)に基づいてハッシュテーブルを構築し、各ハッシュキーを元の行にマッピングします。このハッシュテーブルは、プローブフェーズが開始される前に完全にポピュレートされます。
プローブフェーズ:大きい方のテーブル(lineitem)をスキャンし、各行のキー(l_partkey, l_suppkey)をハッシュします。ビルドテーブルとプローブテーブルの結合キーのハッシュ間のマッチングにより、両方のテーブルの列を組み合わせた出力行が生成されます。マッチしない行はドロップされます(実際には、ここにはハッシュ衝突処理も含まれますが、今は無視します)。
注:左、右、および完全外部結合は、マッチしない行に対する異なるルールで同じビルド/プローブ戦略を使用します。
最低限、このインメモリ結合は3つのテーブルを保持します:ビルドテーブル、プローブテーブル、出力テーブル、そしてビルドサイド上に構築されたハッシュテーブルです。
分散結合
インメモリの場合、すべてのデータは既に同じメモリ空間内に共存しています。テーブルが多くのプロセス/ノード/ランクに分散されている場合、またはテーブルが「ストリーミング」結合のためにバッチ処理されている場合、これはもはや真実ではありません。ランク/プロセスは、レジデントメモリ内の行しか結合できません。最終的にはインメモリ結合が発生しますが、まずビルドテーブルとプローブテーブルの両方のすべてのマッチングキーを同じランクに取得する必要があります。以下の図は、同じ色の様々な行が同じ出力パーティションにシャッフルされ、それらのパーティションが異なるランクに配置される様子を表しています。
分散ハッシュ結合を実行するには、以下の手順を実行する必要があります。
ビルドテーブルをスキャンし、各行の結合キーをハッシュして宛先パーティションを選択します(hash(keys) % n_out_partitions)。各行を所有するランクにパッキングして送信します。
プローブテーブルをスキャンし、同様の方法でルーティングして、プローブキーが同じハッシュを持つビルドキーを既に保持しているランクに着地するようにします。
すべてのランクが送信を完了するまで待ちます。その後、各ランクは、所有するキーに対応する両方のテーブルのすべての行を保持していることが保証されます。
上記を実行したインメモリ結合を、各ランクのローカルスライスで実行します:そのビルド行のハッシュテーブルを構築し、プローブ行でそれをプローブし、マッチを出力します。
最悪の場合、各フェーズが次のフェーズの前に完全にマテリアライズされると、単一のランクは同時に以下のすべてを保持していることになります。
ビルドテーブル(ソーススキャン)
プローブテーブル(ソーススキャン)
ステージングされたビルドテーブル(送信用にパッキング)
ステージングされたプローブテーブル(送信用にパッキング)
シャッフルされたビルドスライス(受信済み)
シャッフルされたプローブスライス(受信済み)
ビルドスライス上のハッシュテーブル
出力テーブル
これが、シャッフリングが計算集約的ではなくメモリ集約的である理由です。ハッシュ計算自体は容易であり、ETLエンジンを構築する際にアウトオブコアシャッフル実装を最優先事項とすべき理由もここにあります。さらに、テーブルをシャッフリングする前に完全にマテリアライズするのではなく、データをストリーミングできることは、メモリ圧力を軽減するために不可欠です。これらの理由から、RapidsMPFはストリーミングアウトオブコアシャッフラーを構築するという当初の目標から始まりました。
RapidsMPF
RapidsMPFは、当初の構想から拡大しました。現在では、2つの大きな部分からなるライブラリです。
1. スピル/アウトオブコアメモリハンドリングのために設計され、高速化されたトランスポートを備えたシャッフルライブラリ
2. ストリーミングデータパイプラインを構築するためのアクターネットワーク
現在、ユーザーはRapidsMPFのシャッフリングコンポーネントのみ(C++またはPythonインターフェース)を採用できます。NeMo-Curatorや、実験的にRay Dataでもこの採用が見られます。最も重要なのは、cuDF Polarsがシャッフルとアクターネットワークの両方にRapidsMPFを使用していることです。フォローアップ投稿でアクターネットワークについて詳しく説明するか、今すぐ興味がある場合はストリーミングエンジンに関するセクションを読むことをお勧めします。
私たちのシャッフル実装は、以下の要件を満たす必要があります。
高速であること
スケーラブルであること
VRAM(GPU)を超えるデータ(アウトオブコア)で動作すること
再利用可能であること
そして、このブログの残りの部分では、メモリ圧力下でのシャッフリングに焦点を当てます。
ベンチマーク設定
cuDF/RapidsMPFには、使いやすいC++ベンチマーク `bench_shuffle` があり、様々なハードウェア(トランスポート、GPU数など)や様々な設定(入力/出力パーティションサイズ、メモリリソースなど)でRapidsMPFのシャッフリング実装がどのように機能するかを研究するのに役立ちます。以下は、現在の `bench_shuffle` テストがユーザーに公開している機能の完全な内訳と、この記事全体で使用する値です。一般的に、このベンチマークは、ランクごとに調整可能な量のランダムな32ビット(4バイト)整数を構築し、すべてのデータをシャッフルして完了します(ここでは結合はありません、シャッフルのみです)。
フラグ | 意味 | ここで使用する値
---|---|---
-C <name> | コミュニケータ | ucxx
-c <n> | 列数 | 10
-r <n> | テスト実行回数 | 10
-w <n> | ウォームアップ実行回数 | 3
-n <n> | ランクあたりの行数 | 536870912 (4バイト/行で列あたり2 GiB)
-p <n> | ランクあたりの入力パーティション数 | 1
-o <n> | ランクあたりの出力パーティション数 | 8 (ランクごとに1つ)
-m <name> | RMMメモリリソースプール | pool
-l <n> | デバイスメモリ制限 (MiB) | 省略 = 無制限 (バイナリデフォルト -1); 32768 から 12288 までスピルスイープ
-s | 出力破棄を有効にする(ストリーミングをシミュレート) | フラグ、常に設定
-x | メモリプロファイリングを有効にする | フラグ、常に設定
-g | 事前パーティション化された入力テーブルを使用する | フラグ、常に設定
すべてのテストで、単一のDGXB200を使用し、`rrun`(NUMAノードにプロセスをバインドできるMPIライクなマルチプロセス起動ツール)を使用してシャッフルを起動します。
シンプルなシャッフリング
DGXB200は8つのBlackwell GPU(それぞれ180GB VRAM)と2つのIntel® Xeon® Platinum 8570プロセッサを備えています。ベースラインを取得するために、まずすべてのGPUに余裕で収まるデータをシャッフルすることから始めます。
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
ここで、ベンチマークを3回ウォームアップし、その後10回実行します。536,870,912行(-n)、10列(-c)、ランクあたり1つの入力パーティション(-p)、そしてデータは8つの出力パーティション(-o)にシャッフルされます。また、UCXX/UCXを使用して、高速化されたトランスポート/GPUDirect RDMAを有効にしています。
536,870,912行 * 4バイト(32ビット整数)= 列あたり2 GiB
10列 * 2 GiB = ランクあたり20 GiB
8ランク * 20 GiB = 合計160 GiB
# ランクあたり20 GiBの例出力
[6:PRINT:0:2026-09-16 02:09:53.934930934] elapsed: 17.91 s | local throughpu