# データフローエンジン ## 定義 データフローエンジンとは、MapReduceの問題点を解消するために開発された分散バッチ処理の実行エンジンで、SparkとFlinkが代表例である。MapReduceが各ジョブを独立したサブジョブ(map/reduce)へ分解して実行するのに対し、データフローエンジンはワークフロー全体を1つのジョブとして扱う。低水準API(1レコードずつ処理するユーザー定義関数の繰り返し呼び出し)に加え、joinやgroup byのような高水準演算子も提供し、演算子はmap/reduceという固定の役割に縛られず柔軟に組み合わせられる。研究システムDryad・Nepheleに端を発する設計であり、内部的にはシャッフルアルゴリズムで結合・集約を実装する。(Source: [[@2026__OReilly__Designing Data-Intensive Applications 2E - Chapter 11 Batch Processing]] "Dataflow Engines") ## MapReduceに対する優位点 - ソートのような高コストな処理は必要な箇所でのみ実行され、map/reduceの間で常に発生するわけではない。 - シャーディングを変えない演算子(mapやfilter等)が連続する場合、単一タスクへ統合してデータコピーのオーバーヘッドを削減できる。 - ジョインやデータ依存関係が明示的に宣言されるため、スケジューラがデータ局所性を意識した配置(生産タスクと消費タスクを同一マシンに置き、共有メモリバッファでデータをやり取りする)を最適化できる。 - 演算子間の中間状態は多くの場合メモリまたはローカルディスクに保持すれば十分で、分散ファイルシステム/オブジェクトストアへの複製書き込みより少ないI/Oで済む。MapReduceはこの最適化をマッパー出力にのみ適用するが、データフローエンジンは全中間状態へ一般化する。 - 演算子開始に新規プロセス起動を必要とするMapReduce(タスクごとに新規JVMを起動)と異なり、既存プロセスを再利用でき起動オーバーヘッドが小さい。 - 演算子は先行ステージの完了を待たずに入力が到着し次第実行を開始できる。 (Source: [[@2026__OReilly__Designing Data-Intensive Applications 2E - Chapter 11 Batch Processing]] "Dataflow Engines") ## 耐障害性の設計差 SparkとFlinkは耐障害性の実現方式が異なる。Sparkは中間データをメモリに保持し(収まらない場合はローカルディスクへ「スピル」)、中間データがどう計算されたか(リネージ)を記録することで、データ喪失時に再計算できる。Flinkは定期的にタスクのスナップショットをチェックポイントする方式を採る。両者ともMapReduceが常に中間データをDFSへ書き戻す方式より効率的である。(Source: [[@2026__OReilly__Designing Data-Intensive Applications 2E - Chapter 11 Batch Processing]] "Handling faults") ## 高水準プログラミングモデル データフローエンジンは低水準APIに加え、SQL・DataFrame API・グラフ処理API(Spark GraphX、Flink Gelly)を提供する。SQLはHive・Trino・Spark・Flinkのコストベースクエリオプティマイザによってジョイン順序の自動最適化を受ける。DataFrame APIはPandas/Rのローカル志向のAPIを分散環境へ拡張したもので、Sparkはメソッド呼び出しをクエリプランへ変換してから最適化・実行する点でPandasの即時実行と異なる。(Source: [[@2026__OReilly__Designing Data-Intensive Applications 2E - Chapter 11 Batch Processing]] "Query Languages", "DataFrames") ## 横断的知見 - **1992年の並列データベース理論が体系化した2つの並列化形態が、34年後には単一プロセス内(DuckDB)とマルチマシン分散処理(Spark/Flink)の両方で独立に再発見されている**: DDIA第11章は、データフローエンジンが「シャーディングを変えない演算子の連結」(パイプライン並列化に相当)と「データを分割した演算子の複製実行」(パーティション並列化に相当)の両方を採用すると説明する。これは[[プッシュ型パイプライン実行]]が記述するDuckDBの単一プロセス内パイプライン実行(スレッドごとの演算子連鎖+ソース分割による並列インスタンス)と、規模もアーキテクチャも異なるにもかかわらず同じ2分類に帰着する。[[並列データベース]]がまとめるDeWitt/Gray(1992)の分類が、分散クラスタと単一マシンのマルチコアという全く異なるスケールの実装に同時に当てはまることが、両ソースの突き合わせで確認できる。(Source: [[@2026__OReilly__Designing Data-Intensive Applications 2E - Chapter 11 Batch Processing]], [[プッシュ型パイプライン実行]], [[並列データベース]]) ## 未解決の問い - データフローエンジンの「データ局所性を意識したスケジューリング」は、具体的にどのようなアルゴリズムで生産タスクと消費タスクの配置を決定するか(本章は概念のみ言及)。 - Spark(リネージ再計算)とFlink(チェックポイント)の耐障害性設計は、性能特性(復旧時間、定常時オーバーヘッド)がどう異なるか。 - DataFrame APIにおける「ローカルDataFrameはインデックス付き・順序ありだが、分散DataFrameは通常インデックスなし・順序なし」という差異は、既存コードの移行時にどのような具体的な性能劣化・挙動差を引き起こすか。 ## 関連 - ソース: [[@2026__OReilly__Designing Data-Intensive Applications 2E - Chapter 11 Batch Processing]] - 概念: [[MapReduce]](前身のアーキテクチャ) / [[シャッフルと分散結合]](内部実装アルゴリズム) / [[プッシュ型パイプライン実行]](単一プロセス版の類似設計) / [[並列データベース]](並列化形態の理論的分類) - 実体: [[Apache Spark]] / [[Apache Flink]] - 書籍: [[Designing Data-Intensive Applications 2nd Edition]] ## 出典 - [[@2026__OReilly__Designing Data-Intensive Applications 2E - Chapter 11 Batch Processing]]("Dataflow Engines"・"Handling faults"・"Query Languages"・"DataFrames" 節)