Rust 実践講座 #5 rayon 並列化 — データ競合のない map-reduce

読了 5分

第 4 回で release ビルドによって 1 コアの性能を使い切りました。しかし最近の機械にはコアが 8 個、16 個とあります。今回はそのコアに集計を分ける話であり、同時に、基礎第 4 回で予告だけしていた一文「借用規則がデータ競合をコンパイル時に遮断する」を実際に目撃する回でもあります。

まず、並列化が得になる条件から #

とにかく分ければ速くなる、というものではありません。判断基準はボトルネックの位置です。

  • IO バウンド(ディスクの読み取りがボトルネック): コアを増やしてもディスクは 1 つです。並列化の利得はほぼありません。
  • CPU バウンド(パース・集計の演算がボトルネック): コア数ぶんの利得の余地があります。

loglens の作業は、行ごとの文字列分解と数値変換というパースが重く、最新の SSD のシーケンシャル読み取りはそれよりずっと速いです。つまり CPU バウンド側で、並列化の候補です。戦略も 1 つ変えます。第 3 回ではメモリを節約するためにストリーミングを選びましたが、並列処理はデータを切って配る必要があるので、ファイルを丸ごと読んでおいて断片をスレッドに配る方が単純で速いです。「基本はストリーミング、並列化が必要なコマンドだけ丸ごと読み」というコマンド別の戦略が、このツールの決定です。

rayon — イテレータを並列に #

インストール
cargo add rayon

rayon の核心は「使い慣れたイテレータチェーンの並列版」であることです。

src/main.rs
use rayon::prelude::*;

fn cmd_stats_parallel(path: &Path) -> anyhow::Result<()> {
    let text = std::fs::read_to_string(path)
        .with_context(|| format!("ログファイルを読み込めません: {}", path.display()))?;

    let stats = text
        .par_lines()                       // 並列の行イテレータ
        .fold(Stats::default, |mut acc, line| {
            acc.feed(line);                // スレッドごとの部分集計
            acc
        })
        .reduce(Stats::default, Stats::merge); // 部分集計の統合

    print_stats(&stats);
    Ok(())
}

lines()par_lines() に変えると、rayon が行をコア数に合わせて分割し、スレッドプールに配り、結果を集めます。スレッドの生成も断片サイズの計算もコードにはありません。構造は map-reduce そのものです。foldスレッドごとに Stats を 1 つずつ作って自分の担当行を集計し、reduce が部分集計を 1 つに統合します。新しく必要なのは統合の関数だけです。

src/main.rs
impl Stats {
    fn merge(mut self, other: Stats) -> Stats {
        self.total += other.total;
        self.parsed += other.parsed;
        self.failed += other.failed;
        self.bytes_sum += other.bytes_sum;
        for (status, count) in other.by_status {
            *self.by_status.entry(status).or_insert(0) += count;
        }
        self
    }
}

数百万行のログなら、コア数に近い倍率を期待できます。逆に数千行のファイルでは、スレッドに配るオーバーヘッドが利得を食いつぶして、かえって遅くなることもあります。第 4 回のストップウォッチで自分のデータで測るのが、いつでも結論です。

なぜ Mutex 共有ではなく fold・reduce なのか #

他の言語の経験者が最初に思いつく設計は、「HashMap を 1 つ作り、ロックで守りながら全員で更新」でしょう。Rust でも可能です(Mutex で包めばよい)。しかしこの設計は行ごとにロックを取る構造なので、スレッドがロックの前に行列を作るロック競合が、並列化の利得をそっくり返上してしまいます。コアを 8 個使ったのに 1 コアより遅い、という結果も珍しくありません。fold・reduce は集計の間の共有がそもそもなく、統合は断片の数だけしか起きません。「共有してロックする」より「分けて合わせる」が並列集計の基本形——これがこの回で持ち帰る設計感覚です。

コンパイラが捕まえるデータ競合 #

この回のハイライトは、実は失敗する場面です。fold・reduce が面倒だからと、外側の HashMap をクロージャから直接更新したらどうなるでしょうか?

src/main.rs
let mut by_status: HashMap<u16, u64> = HashMap::new();
text.par_lines().for_each(|line| {
    if let Ok(entry) = parse_line(line) {
        *by_status.entry(entry.status).or_insert(0) += 1; // コンパイルエラー
    }
});
コンパイルエラー
error[E0596]: cannot borrow `by_status` as mutable, as it is a
              captured variable in a `Fn` closure

複数のスレッドが同時に実行するクロージャが 1 つの値を可変で借りようとした瞬間、基礎第 4 回の借用規則(「可変参照は 1 つだけ」)がそのまま発動します。他の言語ならコンパイルは通り、運の悪い日にカウントがわずかにずれる、再現不可能なバグとして現れたはずのコードです。Rust ではデータ競合はデバッグの対象ではなく、コンパイルエラーの一覧です。所有権と借用という基礎講座の投資が、並列コードで利子になって返ってくる瞬間です。

まとめ #

  • 並列化の判断はボトルネックの位置からです。パース・集計のような CPU バウンドが候補で、IO バウンドはコアを増やしても効きません。
  • 並列のコマンドは丸ごと読みに戦略を変えます。ストリーミング(基本)と丸ごと読み(並列)をコマンド別に使い分けるのがこのツールの決定です。
  • rayon は lines()par_lines() に変えるところから始まります。スレッドごとの fold と統合の reduce という map-reduce が並列集計の基本形です。
  • Mutex でマップ 1 つを共有する設計は、行ごとにロックを取って競合で利得を返上します。「分けて合わせる」が定石です。
  • 共有可変状態のミスは借用規則がコンパイルエラーで捕まえます。データ競合が再現不可能なバグではなくビルド失敗になるのが、Rust の並列化の決定的な利点です。
  • 次回はテストです。パーサーの単体テストから、assert_cmd でツール全体を実行してみる統合テストまで、ここまで作ったものに安全網を張ります。
X