AI・機械学習
10GBのRAMで数十億規模のグラフ上のアルゴリズムを実行: DataFusionに夢中
Algorithms on billion-scale graph using 10GB RAM: I love DataFusion (semyonsinchenko.github.io)
要約
Apache DataFusionを使用して、ディスクへのオフロードを最大限に活用し、ランダムアクセスではなくバルクスキャンに依存するグラフMapReduceを実装しました。これにより、10億エッジのグラフでPageRankを5GBのメモリで、20億エッジのグラフで弱連結成分を10GBのメモリで計算可能になりました。これは、従来のメモリにグラフ全体をロードする必要があるアルゴリズムでは不可能でした。
全文翻訳
TLDR; Apache DataFusionでグラフMapReduceを実装しました。可能な限りすべてをディスクにオフロードし、ランダムアクセスではなくバルクスキャンに依存するようにアルゴリズムを設計しました。DataFusionはスピルオーバー、ソートマージ結合、集計、プランニング、実行を処理するため、私のコードは非常に軽量です。システムd-runを介してハードメモリ制限付きで実行するという厳格なモードでテストしました。動作します。
もちろん、いくつかの問題に遭遇しました。例えば、極端なシナリオではFairSpillPoolからのデッドロックを頻繁に経験し、SMJ(ソートマージ結合)でディスク上のデータの事前ソートを使用する方法をまだ見つけていません。しかし、動作します。Graphalyticsデータセットのgraph500-26(10億エッジの有向グラフ)でPageRankを5GBのメモリで計算できます。あるいは、同じデータセットコレクションのtwitter_mpi(20億エッジのグラフ)で弱連結成分を10GBのメモリで特定できます。NetworkXもIgraphもこれはできません。ほとんどの既存のグラフアルゴリズムは、グラフがメモリに収まることを必要とします。以前は、数十億規模のグラフ分析にはApache SparkとGraphFramesが必要だと思っていました。しかし今では、ラップトップがあれば十分だと思います。Apache DataFusionをグラフ分析に使用することについての私の古い意見は完全に変わりました。
セットアップ
2つのタスクをテストしました。
PageRank
PageRankとは何ですか?
Graphalyticsデータセットのgraph500-26でPageRankを計算するタスクです。
キーバリュー
ノード数: 32,804,978
エッジ数: 1,051,922,853
有向: False
メモリ制限: 5 GB
DataFusion Pool Size: 4 GB
PageRankは最も人気のあるグラフ中心性アルゴリズムの1つであり、検索結果ランキングから不正防止スコアリングまで使用されます。私のDataFusion実装は古典的なPregelです。バルク同期並列アルゴリズム(Map-Reduceとも呼ばれる)を、結合と集計を使用して表現しました。SparkのGraphFramesライブラリのコアにあるものと非常によく似ています。
弱連結成分
弱連結成分とは何ですか?
同じデータセットのtwitter_mpiで、20億エッジのグラフの弱連結成分を特定するタスクです。
キーバリュー
ノード数: 52,579,682
エッジ数: 1,963,263,821
有向: True
メモリ制限: 10 GB
DataFusion Pool Size: 8 GB
WCC(弱連結成分)は、あらゆるID(エンティティ)解決問題の中核部分です。例えば、異なるシステムからのデータをトランジティブIDを通じて重複排除する必要がある場合、WCC問題に直面します。私のDataFusion実装は、「In-database connected component analysis」、Bögeholzら、arXiv 1802.09478に基づいています。SparkのGraphFramesのために同じアルゴリズムをすでに実装していたので、それは明白な選択でした。
結果
PageRank
簡単な部分です。スケーラビリティを証明するためにSMJを使用しましたが、ノードが小さい(32M)ためHJ(ハッシュ結合)を使用することも可能です。PageRankの状態は単純です。1つの列(rank、f64)、1つの列(out-degree、i64)、1つの参加フラグ(bool)です。HJを使用するとより高速です。PageRankは有向エッジ上で動作するため、グラフを対称化する必要はありません。エッジをディスクにオフロードし、収束するまで状態を更新しながら(そしてリネージを断ち切るためにディスクにもオフロードしながら)反復します。
計算時間は長いです。15回の完全な反復で約30分かかります。しかし、セットアップは速度ではなくメモリに関するものです。10億エッジのグラフ分析に、より現実的な数値を設定すれば、十分に高速に動作するでしょう(テスト済み)。グラウンドトゥルースと数値を比較しました。100%一致(0.0001の許容誤差)。ここでも多くの最適化が可能です。理論的には、エッジを範囲でバケット化したり、範囲パーティショニングのようなものを実行したりすることが可能であり、SMJが各反復でトリプレットを取得するために最大の結合側(エッジ)を再度ソートする必要がなくなります。また、Parquetが最良の選択肢であるかどうかは100%確信が持てません。結合と集計を融合することも興味深いでしょう。各Pregel反復は、edges <-[join] nodes-state -> group by + agg -> [join] -> nodes-state -> update nodes-state のようなものです。最初の2つのステージを融合できれば、パフォーマンスの観点から大きな勝利となる可能性があります。一方、DataFusionでこれをどのように行うかはまだわかりません。学ぶべきことはたくさんあります。
WCC
最も難しい部分です。20億エッジのTwitterグラフはすでに巨大です(CSV形式ではエッジだけで30GB!!!)。しかし、WCCではエッジを対称化する必要があります(またはsrc、dst、およびdstをsrc、srcをdstとする結合+その上のdistinctを実行する)。そのため、ピーク時には8GBのDataFusionプールのみを使用して約40億エッジを処理しています。フローが最初の数回の反復を生き残ると、収縮プロセスによってエッジの量が劇的に減少し、アルゴリズムは10分で低メモリ負荷で終了します。
sem@fedora:~/github/graphframes-rs$ systemd-run --user --scope \
-p MemoryMax=10G -p MemorySwapMax=0 \
-p AllowedCPUs=0-1 \
--setenv=RUST_LOG=graphframes_rs=info,datafusion=warn \
./target/release/run-algorithm twitter_mpi-v.parquet twitter_mpi-e.parquet wcc 42 file:///var/home/sem/Downloads/gf_wcc_out 8G 2
実行中のユニット: run-p316509-i284528.scope; 呼び出しID: 742f9296d31d426580b7ec8213422cf1
[2026-07-05T05:37:21Z INFO graphframes_rs::algorithm::connectivity::connected_components] start WCC with run-id 017c0a23-2b20-4ffa-ac6b-6e2cb8d7203e
[2026-07-05T05:52:21Z INFO graphframes_rs::algorithm::connectivity::connected_components] after preparation graph has 3228212374 edges
[2026-07-05T06:13:21Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 1, edges remaining: 840238268
[2026-07-05T06:17:39Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 2, edges remaining: 77322906
[2026-07-05T06:17:57Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 3, edges remaining: 5624128
[2026-07-05T06:17:59Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 4, edges remaining: 1075998
[2026-07-05T06:17:59Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 5, edges remaining: 230838
[2026-07-05T06:17:59Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 6, edges remaining: 97940
[2026-07-05T06:17:59Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 7, edges remaining: 42352
[2026-07-05T06:17:59Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 8, edges remaining: 16720
[2026-07-05T06:17:59Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 9, edges remaining: 8238
[2026-07-05T06:17:59Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 10, edges remaining: 3860
[2026-07-05T06:17:59Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 11, edges remaining: 1488
[2026-07-05T06:17:59Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 12, edges remaining: 982
[2026-07-05T06:17:59Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 13, edges remaining: 514
[2026-07-05T06:17:59Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 14, edges remaining: 132
[2026-07-05T06:17:59Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 15, edges remaining: 120
[2026-07-05T06:17:59Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 16, edges remaining: 40
[2026-07-05T06:17:59Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 17, edges remaining: 18
[2026-07-05T06:17:59Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 18, edges remaining: 10
[2026-07-05T06:17:59Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 19, edges remaining: 6
[2026-07-05T06:17:59Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 20, edges remaining: 4
[2026-07-05T06:17:59Z INFO graphframes_rs::algorithm::connectivity::connected_components] cc forward iteration 21, edges remaining: 2
[2026-07-05T06:17:59Z INFO graphframes_rs::algorithm::conn