MapReduceによる大規模分散データ処理
1. 概要
A. 定義
MapReduceとは、大規模な入力をキー・バリュー(key-value)レコードに分割し、複数ノードで
Map演算を並列実行し、同じ中間キーの値を集約してReduce演算で結果を生成する分散処理プログラミングモデル兼実行フレームワークである。
開発者は主にmap関数とreduce関数で業務計算を表現し、実行系が入力分割、スケジューリング、ネットワーク転送、ソート、障害回復を担当する。これによって業務ロジックとクラスタ運用の複雑さを分離できる。
このモデルはGoogleが2004年に発表した大規模データ処理の論文で体系化され、オープンソースのHadoop MapReduceによって広く普及した。Spark、分散SQL、ストリーム処理エンジンが反復計算や低遅延処理に適する場合でも、分割・シャッフル・集約という原理はデータ基盤に残っている。
B. 背景と必要性
単一サーバのバッチ処理は、データがメモリ・ディスク・CPU能力を超えるとより大きな機器への垂直拡張に依存する。ログ、クリックストリーム、センサ、取引記録は急増するため、単一機器だけでは費用と障害リスクを抑えにくい。
分散処理はデータを複数サーバに分け、計算を並列化する。しかし直接分散プログラムを書くと、分割、障害、再試行、結果結合、通信制御をアプリケーションが負担する。MapReduceはこれらの共通処理をフレームワークに移す。
データ局所性(data locality)も重要である。すべての入力を中央サーバへ移すとネットワークがボトルネックになるため、MapReduceは入力ブロックを保有するノードの近くにmapタスクを配置する。
C. 主要目標
目的は単にサーバを増やすことではない。スループット拡張性、障害耐性、プログラミングの単純性、データ局所性、運用自動化を同時に実現することである。
入力を独立した分割として並列化し、同じキーを再分配して集約し、失敗したタスクだけを再実行する。一方、中間結果をディスクへ保存するため、対話型・反復型処理では別のエンジンが適する場合がある。
2. プログラミングモデルと実行構造
A. キー・バリュー抽象化
MapReduceは入力と出力を<key, value>対として扱う。入力キーはファイルオフセットやDBキー、値はテキスト行・JSON文書・センサレコードになりうる。
map関数は入力レコードを読み、0個以上の中間対を出力する。reduce関数は一つの中間キーと、そのキーに属する値のリストを受け取る。
map(k1, v1) -> list(k2, v2)
reduce(k2, list(v2)) -> list(k3, v3)
mapはレコード単位の変換・フィルタ、reduceは同じ論理キーの集約・結合を担当する。この境界が独立処理とグループ処理を明確にする。
B. 全体アーキテクチャ
flowchart LR
IN[入力ファイル・テーブル] --> SPLIT[Input Split]
SPLIT --> MR[Mapper]
MR --> COMB[任意のCombiner]
COMB --> PART[Partitioner]
PART --> SHUF[Shuffle・Sort]
SHUF --> RED[Reducer]
RED --> OUT[分散ファイルシステム出力]
RM[ResourceManager] -.スケジュール・資源.-> MR
RM -.スケジュール・資源.-> RED
NM[NodeManager] -.実行・状態報告.-> MR
NM -.実行・状態報告.-> RED
HadoopではResourceManagerが資源と配置を調整し、NodeManagerがコンテナ実行と状態報告を担当する。ApplicationMasterはmap・reduceタスクを追跡し、失敗タスクの再試行を要求する。
InputFormatとRecordReaderがファイルを論理レコードへ変換する。テキストなら行を値、オフセットをキーにできるが、レコード境界を保つことが正確性の前提になる。
mapperは中間出力をローカルにバッファし、パーティション別・キー順に保存する。これはジョブ成功時に確定するまで一時出力として扱われる。
C. 処理ライフサイクル
- クライアントが入力・出力、mapper、reducer、パーティション数を送信する。
- フレームワークが入力をsplitへ分割する。
- 可能ならデータを保持するノードへmapタスクを配置する。
- mapperがレコードを読み、中間キー・バリューを生成する。
- 任意のcombinerがローカル部分集約を行う。
- partitionerが各キーを担当reducerへ割り当てる。
- reducerがmapperの結果を取得し、マージソートする。
- reducerがキー別に値をグループ化して計算する。
- 出力フォーマットが結果を保存し、ジョブをコミットする。
Shuffleは単なるコピーではない。すべてのmapperからreducerの担当範囲を取得し、キー順にマージするため、時間・I/Oを支配する段階になりやすい。
sequenceDiagram
participant C as クライアント
participant AM as ApplicationMaster
participant M as Mapperノード
participant R as Reducerノード
C->>AM: ジョブ送信(入力・出力・関数)
AM->>M: split別map配置
M->>M: 中間対生成・ローカルソート
M->>R: partition別shuffle
AM->>R: reduce起動
R->>R: キー別グループ化・集約
R-->>C: 出力コミット・状態報告
ApplicationMasterはレコードを中央処理しない。配置と状態を追跡し、再試行を要求する制御面であり、実際の計算はmapperとreducerが行う。
3. MapとReduce構成要素の原理
A. Mapper
mapperは通常、各レコードを独立に処理する。例えばログのステータスが500なら<service, 1>を出力し、不正レコードを早期に捨ててshuffle量を減らす。
入力型・中間型の変換、データ清浄、パーティションキー設計もmapperの責任である。キー設計が悪いと一つの業務グループが分裂し、または人気キーが一つのreducerを過負荷にする。
B. Combiner
combinerはmapperとreducerの間で行う任意のローカル集約である。WordCountなら、同一mapper内のcloud 10,000件を<cloud, 10000>へ圧縮できる。
部分結果を再び結合しても結果が同じになる必要がある。合計・個数・最小・最大は適するが、平均は合計と個数を一緒に伝える必要がある。
combinerの実行有無や回数は保証されない。したがって正確性をcombinerだけに置いてはならない。
C. Partitioner
partitionerは中間キーを担当reducerへ割り当てる。ハッシュ方式が一般的だが、日付、地域、顧客群など業務範囲で分けるなら独自partitionerを使える。
同じ論理キーは同じreducerへ送られ、全体の負荷は均等でなければならない。人気商品・地域・日付が一つへ集中するとデータスキューが発生する。
hot keyにsaltを付けて複数partitionへ分け、二段目で再集約する方法がある。ただし追加の結合段階が必要なため、正確性と費用を併せて評価する。
D. ShuffleとSort
shuffleはmapperのpartitionをreducerへ転送する。ネットワークとディスクI/Oを消費し、ジョブ遅延の主因になりやすい。
sort・mergeは同じキーを連続させるため、reducerは全データをメモリへ載せずキーグループをストリーム処理できる。
早期フィルタ、適切なcombiner、中間圧縮、partition数の調整が基本的な最適化である。ただし圧縮CPUが通信削減効果を上回る場合もある。
E. Reducer
reducerは一つのキーに属する値を集約、整列、結合、重複排除、統計計算する。入力がキー順なのでキーが変わる時点で状態を解放できる。
reducerは再実行される可能性がある。外部書き込みには冪等キー、一時出力、原子的コミットを使い、重複API呼出しを防ぐ。
4. 障害耐性と性能設計
A. タスク再実行
大規模クラスタではディスク故障、ネットワーク断、プロセス停止、ノード過負荷が起きる。MapReduceはジョブ全体ではなく失敗タスクを別ノードで再実行する。
入力ブロックの複製により別コピーからmapを再実行できる。stragglerにはspeculative executionで同じ処理を別ノードでも実行し、先に成功した結果を採用できる。
外部決済など非冪等副作用は重複実行で危険になるため、純粋関数または冪等な出力プロトコルが必要である。
B. 局所性とファイル設計
性能はCPUよりデータ移動に制限される場合が多い。ローカルブロックを読めば通信を避けられるが、リモート配置では転送費用が増す。
小ファイルが多いとsplitとタスク起動の負荷が増え、splitが少なすぎると並列性が下がる。ファイル、ブロック、タスク、ノードの数を一緒に調整する。
C. コストと観測
map計算量を(M)、中間レコード数を(I)、reducer数を(R)とすれば、総費用はmap、書込み、(I)のshuffle・sort、reduceの合計として考えられる。
監視ではmap入出力、combiner削減率、shuffleバイト、待ち時間、reducer間の入力差、失敗・再試行数を見る。特定reducerだけ遅ければキー偏りを疑う。
繰返しジョブでは平均だけでなくp95・p99遅延と入力量の基準線を管理する。
5. 例と応用
A. WordCount
mapperが<word, 1>を出し、combinerが局所集約し、reducerが部分合計を足して<word, total>を出力する。異なるノードの同じキーも一つの論理グループになる。
B. 転置索引
mapperが<word, documentId>を出し、reducerが文書IDを整列・重複排除してposting listを作る。人気語に値が集中するため、stop word除去や追加マージが必要になることがある。
C. ログ・取引集計
日別サービス別エラー数なら<date|service, 1>、顧客別合計なら<customerId, amount>を作る。規制データでは暗号化、マスキング、権限、監査、保存期間を別途設計する。
6. 代替技術との比較
MapReduceはディスク型バッチと自動再試行に強いが、反復処理では中間保存が遅延を増やす。データ量だけでなく遅延、状態、再利用、SLAで選定する。
| 区分 | MapReduce | Spark | 分散SQL | ストリーム処理 |
|---|---|---|---|---|
| 基本 | バッチ段階 | キャッシュDAG・バッチ | 宣言的問合せ | 連続イベント |
| 中間データ | 主にディスク | キャッシュ・メモリ | エンジン計画 | 状態・チェックポイント |
| 強み | 単純性・再試行 | 反復・低遅延 | SQL生産性 | リアルタイム |
| 弱み | shuffle・ディスク遅延 | メモリ調整 | UDF制約 | 状態・重複管理 |
SparkはRDD、DataFrame、DAGでmap/reduceの考えを拡張し、再利用データをキャッシュできる。しかしメモリ圧迫やshuffleコストが消えるわけではない。
分散SQLはjoin、filter pushdown、partition pruningを自動最適化できる。非定型レコードや外部処理の組合せではMapReduceが扱いやすい場合もある。
ストリーム処理は無限入力をウィンドウと状態で処理し、イベント時刻、watermark、遅延到着、重複、処理保証を別途定義する。
7. 深掘り — 現代データ基盤での位置と出題
MapReduceは製品名よりも分散データフローの考え方として理解するとよい。入力分割、キー再配分、部分集約、ソート、最終結合は分散SQLにも現れる。
オブジェクトストレージとコンテナでは保存と計算が分離するため、データ局所性はローカルディスクからファイル形式、partition pruning、キャッシュ、通信費最適化へ移る。
技術士論文ではMap → Combine → Partition → Shuffle/Sort → Reduceを図示し、WordCountやログ集計でキー・バリューのグループ化を説明する。その後、局所性、再試行、スキュー、shuffleボトルネックを運用へ接続するとよい。
8. 考慮事項と示唆
- キー・バリュー集約に適するか確認する。 多段joinで中間データが膨張するならSQL、グラフ、ストリーム技術を検討する。
- shuffleを費用中心として管理する。 中間量、reducer差、ネットワーク、ディスクを測定してからクラスタを拡大する。
- combinerは任意実行と考える。 代数的に安全な部分集約だけを置き、正確性はreducerで保証する。
- 再実行を安全にする。 冪等キー、一時パス、原子的コミットで外部副作用の重複を防止する。
- スキューを業務キーで解決する。 reducer数だけではhot keyの集中は解消しない。
- ファイルとsplitを同時設計する。 小ファイルをまとめ、形式とpartition pruningによる実読込量を確認する。
- ガバナンスを実行経路に含める。 複製・中間データの暗号化、アクセス制御、個人情報マスキング、監査を行う。
- SLAと費用で選ぶ。 MapReduce、Spark、SQL、ストリームを実行時間、復旧時間、費用、運用難度で比較する。
参考資料
- Dean, J. and Ghemawat, S., “MapReduce: Simplified Data Processing on Large Clusters”, Google Research. https://research.google/pubs/mapreduce-simplified-data-processing-on-large-clusters/
- Apache Hadoop, “MapReduce Tutorial”. https://hadoop.apache.org/docs/current/hadoop-mapreduce-client/hadoop-mapreduce-client-core/MapReduceTutorial.html
- Apache Hadoop, “HDFS Architecture Guide”. https://hadoop.apache.org/docs/current/hadoop-project-dist/hadoop-hdfs/HdfsDesign.html
一言まとめ: MapReduceはキー・バリュー入力を分割し、Map・Combine・Shuffle/Sort・Reduceで並列処理することで、データ局所性、タスク再試行、キー集約による大規模バッチの拡張性と障害耐性を実現する。