← 一覧へ
AI・データ
#MapReduce#Hadoop#분산처리#Map#Reduce#Shuffle#데이터지역성#빅데이터배치
最終更新 · 2026-09-28

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. 処理ライフサイクル

  1. クライアントが入力・出力、mapper、reducer、パーティション数を送信する。
  2. フレームワークが入力をsplitへ分割する。
  3. 可能ならデータを保持するノードへmapタスクを配置する。
  4. mapperがレコードを読み、中間キー・バリューを生成する。
  5. 任意のcombinerがローカル部分集約を行う。
  6. partitionerが各キーを担当reducerへ割り当てる。
  7. reducerがmapperの結果を取得し、マージソートする。
  8. reducerがキー別に値をグループ化して計算する。
  9. 出力フォーマットが結果を保存し、ジョブをコミットする。

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、ストリームを実行時間、復旧時間、費用、運用難度で比較する。

参考資料


一言まとめ: MapReduceはキー・バリュー入力を分割し、Map・Combine・Shuffle/Sort・Reduceで並列処理することで、データ局所性、タスク再試行、キー集約による大規模バッチの拡張性と障害耐性を実現する。