
2026/08/01 0:53
10GB のメモリで数十億規模のグラフ処理するアルゴリズム:DataFusion が大好きです
RSS: https://news.ycombinator.com/rss
要約▶
Japanese Translation:
Apache DataFusion は、標準的なノートパソコン上でバイリオン規模のグラフ解析を遂行する能力を証明し、Apache Spark といった高価なクラスターインフラストラクチャの必要性を取り除きました。SQL スタイルのジョインとアグリゲーションを通じて効率的な大規模並列同期(Map-Reduce)アルゴリズムを実装し、処理をディスクにオフロードすることで、同システムはわずか 5GB のメモリを用いて、約 30 分以内に PageRank を計算しました。その対象は 15 回のフル反復にわたるバイリオンエッジの有向グラフ(Graphalytics
graph500-26、エッジ数 1.05B、ノード数 3280 万)であり、計算結果は真値と完全に一致しました。同様に、同システムは Bögeholz et al. の「インデータベース連結成分解析」に基づき、約 36 分以内に 22 回の前方反復を用いて、Twitter データセット(twitter_mpi、エッジ数~19.6B、ノード数 5260 万)の弱連結成分を特定しました。その際、メモリプールは 8GB が使用されました。これらの結果は、systemd-run を介した厳格なメモリ制限などの条件下でも検証されています。依然として、スパイルプールのディスク前ソート化の欠如による極限メモリシナリオにおける稀なデッドロックなど、若干の課題は残っていますが、この飛躍的な進歩により、研究者や小規模チームが安価なハードウェア上で地元の環境で高度な解析を実行するための障壁が大幅に低下しました。本文
Apache DataFusion を用いた大規模グラフ Map-Reduce 実装報告
概要 (TL;DR)
- アプローチ: Apache DataFusion を基盤に、大規模グラフに対する Map-Reduce アルゴリズムを実装しました。
- 設計方針:
- 処理をディスクへオフロードし、ランダムアクセスではなくバッチスキャンを基盤としたアルゴリズムを採用しています。
- DataFusion がスパイラ(Spill)、ソートマージ結合(SMJ)、集計、プランニング、実行を一括で担当するため、実装コードは極めて軽量です。
- 検証環境: 厳密なモードでテストを行い、
を用いてハードメモリ制限下での動作確認を実施しました。systemd-run
主要な成果と課題
- PageRank (10 億エッジ):
- データセット
(3,280 万ノード、10.5 億エッジ)に対し、5GB のメモリで計算可能です。graph500-26
- データセット
- 弱連結成分 (WCC) (20 億エッジ):
- データセット
(5,257 万ノード、19.6 億エッジ)に対し、10GB のメモリで全弱連結成分を同定可能です。twitter_mpi
- データセット
- 比較優位:
- NetworkX や Igraph ではこれらが不可能でしたが、現在はノート PC 一台でも対応可能になりました。
- 従来の「Apache Spark と GraphFrames が必須」という認識から、DataFusion の大規模グラフ解析への活用姿勢を根本的に変えました。
- 現時点での課題:
- FairSpillPool: 極限のシナリオではデッドロックが頻発します。
- SMJ 実装: ディスク上で事前ソートされたデータを使用する SMJ の実装は未対応であり、現状は動作していません(※)。
設定 (Setup)
2 つの主要タスクについてテストを行いました。
PageRank
- 概要: Graphalytics データセット
に対する PageRank 計算です。検索結果のランキング付けや不正行為検知などに利用されます。graph500-26 - アルゴリズム:
- 古典的な Pregel(Bulk Synchronous Parallel: Map-Reduce)を採用しています。
- DataFusion の Join と Agg を用いて表現し、Spark の GraphFrames ライブラリのコア部分と類似した構造です。
使用データセット詳細
| 項目 | 値 |
|---|---|
| ノード数 | 32,804,978 |
| エッジ数 | 1,051,922,853 |
| 有向 (Directed) | False(非有向) |
| メモリ制限 (Memory Limit) | 5 GB |
| DataFusion プールサイズ | 4 GB |
弱連結成分 (Weakly Connected Components, WCC)
- 概要: Graphalytics データセット
に対する全弱連結成分の同定です。異なるシステム間の ID 重複を除去する「アイデンティティ解決」の中核アルゴリズムです。twitter_mpi - アルゴリズム: 「データベース内の連結成分分析」(Bögeholz 他,arXiv 1802.09478)に基づく実装です(Spark GraphFrames でも同様のアルゴリズムが実装済み)。
使用データセット詳細
| 項目 | 値 |
|---|---|
| ノード数 | 52,579,682 |
| エッジ数 | 1,963,263,821 |
| 有向 (Directed) | True(有向) |
| メモリ制限 (Memory Limit) | 10 GB |
| DataFusion プールサイズ | 8 GB |
結果 (Results)
PageRank
- 特性:
- 比較的単純なタスクですが、スケーラビリティを証明するため SMJ(ソートマージ結合)を利用しました。
- 頂点サイズが小さいため、HJ(ハッシュ結合)でも可能です。
- 実装詳細:
- エッジと状態をディスクへオフロードし、イテレーションを繰り返します(ラインジンを破断させることなく処理)。
- 収束するまで繰り返し実行します。
- 有向エッジのみを対象とするため、グラフの非対称化は不要です。
- 計算時間:
- 約30 分(15 フルイテレーション)を要しますが、ボトルネックは処理速度ではなくメモリ設定です。
- 現実的な規模(10 億エッジ級)であれば十分高速に動作します(テスト済み)。
- 精度確認:
- グラウンド・トゥリス(正解値)との照合結果:完全一致(許容誤差 0.0001)。
将来の最適化可能性
- 結合効率: SMJ は各イテレーションで最も結合側のデータ(エッジ)を再ソートする必要があり、パフォーマンス向上の余地があります。
- 改善策: エッジを範囲別バケット化やレンジパーティショニングを行うことで回避可能です。
- ストレージ形式: Parquet が最適な形式か検討が必要ですが、Parquet 以外の選択肢も視野に入れます。
- 結合と集計のフューズ (Join+Agg Fuse):
- 各 Pregel イテレーションは「エッジ Join ノード状態 -> Group By + Agg -> Join -> 更新」という流れです。
- 最初の 2 ステージ(Join + Agg)をフューズできれば、パフォーマンスの劇的向上が見込めますが、DataFusion での実装方法についてはさらに学習が必要です。
Weakly Connected Components (WCC)
- 特性:
- 最も困難なパートです。Twitter グラフ(20 億エッジ、CSV 約 30GB)に加え、WCC ではエッジを対称化(Union)する必要があります。
- ピーク時には約40 億のエッジを、プールサイズである8GBという限られたメモリで処理します。
- プロセス:
- 最初の数回のイテレーションを乗り越えると、収縮プロセスによりエッジ数が劇的に減少し、メモリ圧力が低下する段階でアルゴリズムは完了します。
- 完了までの所要時間:約 10 分。
実行ログ(一部)
以下は
systemd-run を用いたハードメモリ制限下での実行ログです。
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 Running as unit: run-p316509-i284528.scope; invocation ID: 742f9296d31d426580b7ec8213422cf9 [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:18:23Z INFO graphframes_rs::algorithm::connectivity::connected_components] connected components written to file:///var/home/sem/Downloads/gf_wcc_out after 22 forward iterations num-iterations: 22
精度確認
- Graphalytics: グラウンド・トゥリス(正解データ)を提供しており、検証が容易です。結果は正しいことを確認しました。
メモリ使用状況と集計 (CSV リード)
memory D SELECT column1, count(*) as cnt FROM read_csv('twitter_mpi-WCC', delim=' ') GROUP BY column1 ORDER BY cnt DESC LIMIT 5; ┌──────────┬──────────┐ │ column1 │ cnt │ │ int64 │ int64 │ ├──────────┼──────────┤ │ 1 │ 52,515,193│ │ 27052874 │ 67 │ │ 47269046 │ 44 │ │ 45352761 │ 33 │ │ 17516773 │ 30 │ └──────────┴──────────┘
結果テーブルからの集計
memory D SELECT component, count(*) as cnt FROM results GROUP BY component ORDER BY cnt DESC LIMIT 5; ┌───────────┬──────────┐ │ component │ cnt │ │ int64 │ int64 │ ├───────────┼──────────┤ │ 1 │ 52,515,193│ │ 27052874 │ 67 │ │ 47269046 │ 44 │ │ 45352761 │ 33 │ │ 17516773 │ 30 │ └───────────┴──────────┘