🍰 まずはやさしく
大量のデータを分担して処理する仕組みです。
安価なPCをたくさんつないで計算します。
スマホの膨大な利用記録などを処理します。
この章では仕組みの核となる3つの機能を紹介します。
大規模データを 「複数の安価なPCに分けて並列処理」 するためのオープンソース基盤
🍰 まずはやさしく
昔のビッグデータ処理の主役です。
とても大きなデータを扱うときに使います。
SNSの投稿ログなどの解析に向いています。
ここでは使うべき場面とそうでない場面を学びます。
「ビッグデータ=Hadoop」 と言われた 2010 年代前半の主役。 統計・データ解析の文脈では、 SSDSE のような 数 MB の CSV を 1 台の pandas で扱えるサイズ では Hadoop は 明らかに過剰 ですが、 SNS ログ、 アクセスログ、 IoT センサー値など TB~PB 級になると本領を発揮します。 本ページでは「いつ使うのか/いつ使うべきでないのか」の判断軸を、 SSDSE-B-2026 を題材に説明します。
🍰 まずはやさしく
みんなで分担して作業するチームのようなものです。
計算時間を短くするために使います。
部活の大量のプリントをみんなで分ける感覚です。
ここでは効率よく計算する考え方を解説します。
1TB のアクセスログを集計するシナリオを想像してください。
これが 水平スケーリング(scale-out)。 1 台を高性能化(scale-up)しても CPU・I/O に物理限界があるが、 台数を増やせば理論上どこまでも速くなります(ただし通信オーバヘッドで頭打ち。 これが下の Amdahl の法則)。
Hadoop の発想の キモは「計算をデータのある場所に送る」。 普通は CPU の所にデータを引き寄せますが、 1 TB を毎回引き寄せると LAN が詰まる。 逆にプログラム (数 KB)をデータのある DataNode に配って、 そこで実行させる — これが「データローカリティ」です。
もう 1 つの肝が 「故障は当たり前」 の発想。 1000 台のクラスタなら毎日数台が壊れます。 だから HDFS は同じブロックを 3 つの別ノードにコピーし、 1 台死んでも残り 2 つで読める。 高価な RAID サーバではなく、 安いコモディティ PC を大量に使う思想です。
🍰 まずはやさしく
データを分けて保存し処理する枠組みのことです。
計算の速度や保存量を正確に管理するために使います。
買い物サイトの膨大な注文データを扱うイメージです。
ここでは計算速度や容量に関する数式を学びます。
Hadoop(Apache Hadoop):大規模データを商用クラスタで分散保存・並列処理するオープンソースフレームワークの総称
Hadoop / Spark のクラスタ規模を決めるとき、 「ノードを 2 倍にすれば 2 倍速くなるか?」という素朴な期待は Amdahl 法則 によって裏切られます。 並列化できない逐次部分(シャッフル、 マスタノードでの集約、 ディスク I/O 待ち)が必ず存在するため、 ノード追加の効果は頭打ちになります。
プログラム全体の実行時間を 1 とし、 そのうち並列化可能な割合を p、 逐次実行が必要な割合を (1-p) とします。 ノード数(並列度)を N としたとき、 スピードアップ S(N) は次式で与えられます:
$$ S(N) = \frac{1}{(1-p) + \frac{p}{N}} $$分母の (1-p) は 逐次部分、 p/N は 並列化された部分が N ノードで分担されたあとの実行時間。 N を無限大にしても、 逐次部分 (1-p) が残るため、 スピードアップは 1/(1-p) で頭打ちになります。 これが Amdahl 限界。
| 記号 | 意味 | Hadoop 上での具体例 |
|---|---|---|
| p | 並列化可能な処理の割合 | Map フェーズ(各ブロックを独立に処理) |
| 1-p | 逐次実行が必要な割合 | Reduce 直前のシャッフル、 マスタ集約、 ジョブスケジュール |
| N | 並列度(ワーカノード数) | DataNode 数(典型例: 3〜100 台) |
| S(N) | スピードアップ比 | 1 ノード時に対する高速化倍率 |
Hadoop ジョブの典型値として p=0.9(Map が支配的)を想定すると、 N を増やしたときの S(N) は次のようになります:
| N(ノード数) | S(N)(スピードアップ) | 追加効率 |
|---|---|---|
| 1 | 1.00x | 基準 |
| 2 | 1.82x | +82% (理想 100%) |
| 4 | 3.08x | +208% (理想 300%) |
| 10 | 5.26x | +426% (理想 900%) |
| 100 | 9.17x | +817% (理想 9900%) |
| ∞ | 10.00x | 上限 = 1/(1-0.9) |
⚠️ 10 ノードで 5.26 倍にしかならない。 「ノード追加すれば線形に速くなる」は幻想であり、 SSDSE のような中小データで Hadoop を導入しても費用対効果が低い理由がここにあります。
Amdahl の法則 $S(p) = 1/((1-f) + f/p)$ は、 並列化の限界を示す Hadoop 設計の根本です。 言葉に置き換えると、 「ある処理のうち、 並列化できる部分の比率が $f$、 残り(シリアル部分)が $(1-f)$ の場合、 $p$ ノード使ったときの速度向上は理論上 $1/((1-f)+f/p)$ 倍が上限」となります。 つまり、 仮に $f=0.99$(99% が並列化可能、 1% がシリアル)の理想的なケースでも、 ノード数を無限大に増やすと $\lim_{p\to\infty} S(p) = 1/(1-f) = 100$ 倍で頭打ち。 「1% のシリアル部分があるだけで、 1000 台投入しても 100 倍までしか速くならない」 という重い事実が読み取れます。
これが Hadoop 設計者が 「Shuffle フェーズの最小化」 や 「Combiner(マップ側集約)」 に執着する理由です。 Shuffle は全ノード間通信を伴うのでシリアル部分(あるいは並列度の低い部分)に近く、 ここで詰まると Amdahl の壁にぶつかる。 たとえばワードカウントで Combiner を使えば各 Mapper 内で部分集計してから Shuffle するので転送量が劇的に減ります(10 億単語のテキストでも、 1000 のユニーク単語に集約してから Reducer に渡せる)。
HDFS のストレージコスト $D \times R + M$ は 「ディスクは安いが無限ではない」を表す式。 レプリカ数 $R=3$ は標準値(Google 論文由来)ですが、 これにより 使用容量は実データの 3 倍 になります。 1 PB のログ を持ちたければ 3 PB のディスクが必要。 これがクラウド時代に イレイジャーコーディング(パリティ符号化で 1.5 倍まで圧縮)に置き換わった理由です。 メタデータ $M$ も無視できず、 100 万ファイル × 1 ブロック × 150 byte = 150 MB の NameNode メモリを消費(小ファイル問題の本質)。
MapReduce の計算時間 $T_{total} = T_{map}(n/p) + T_{shuffle}(n) + T_{reduce}(k/p)$ を読むと、 3 段階で時間がかかる場所が見える。 Map と Reduce は $p$(ノード数)で割れる(並列)が、 真ん中の Shuffle は $n$(全データ量)に比例 — ここがネットワーク帯域に支配される。 たとえば $n=1$ TB、 $p=100$ なら Map と Reduce は約 1/100 に短縮されるが、 Shuffle は全ノード間で 1 TB 級の通信が走ります。 10 Gbps LAN でも 1 TB の転送は約 800 秒。 これが「ジョインが遅い」「GROUP BY 多用は危険」の正体です。 Spark の RDD lineage や DataFrame の Catalyst 最適化が、 この Shuffle 削減を狙う技術として登場した文脈につながります。
ブロック数 $N_{blocks} = \lceil D/B \rceil$ は NameNode のメモリ設計に直結。 既定の 128 MB ブロックは、 1990 年代の HDD シーケンシャル読み出し速度(≈100 MB/s)で約 1 秒で読める量として設計された値です。 もし小さくしすぎる(例:64 KB)と $N$ が膨大になり NameNode メモリが破綻、 大きすぎる(例:4 GB)と並列度が下がる(1000 台クラスタで 100 ブロックしか作れないと活かしきれない)。 $B$ の選択は「並列度 vs メタデータ量」のトレードオフで、 ログ集計のような大規模順次読みでは 256 MB ~ 1 GB に設定するのが定石です。 SSDSE-B-2026 のような数 MB の CSV はそもそも 1 ブロックで収まり、 Hadoop の利点が出ない領域だと、 この式から即座に判断できます。
| 記号/用語 | 意味と具体例 |
|---|---|
| HDFS | Hadoop Distributed File System。 既定 128 MB ブロック、 3 レプリカ。 書込み 1 回・読込み多数(WORM)に最適化。 |
| NameNode | メタデータ(どのファイルがどのブロック・どこに)を管理するマスター。 単一障害点だったが、 HA 構成(Active/Standby)で改善。 |
| DataNode | 実データブロックを保持するワーカ。 数百〜数千台で構成。 ハートビートで NameNode に生存報告。 |
| MapReduce | Map(行ごとに変換)→ Shuffle(キーで集約・転送)→ Reduce(集約値を計算)の処理モデル。 Google 2004 年論文が原典。 |
| YARN | Yet Another Resource Negotiator。 ジョブのリソース要求とノード割当を管理するスケジューラ。 MapReduce/Spark/Tez を載せる土台。 |
| Hive | SQL ライクなクエリ言語(HiveQL)で MapReduce/Tez ジョブを生成。 SQL ユーザが Hadoop を使うブリッジ。 |
| Combiner | Mapper 側で部分集計し Shuffle 転送量を削る最適化。 集計が結合則・可換則を満たす場合のみ使える(和、 最大、 集合和など)。 |
| Partitioner | Mapper 出力をどの Reducer に振るかを決める関数(既定はキーの hash mod R)。 偏りがあると一部 Reducer に集中(skew)。 |
/user/me/log.csv を書きたい」と問い合わせポイント:NameNode は メタデータの仲介のみで、 実データは DataNode とクライアントが直接通信。 これが「NameNode はネックにならない」設計の核心。
| テクニック | 何を改善 | 具体例 |
|---|---|---|
| Combiner | Shuffle 量削減 | 和・最大・最小・集合和など結合則を満たす操作。 WordCount で 1000 倍削減可 |
| In-Mapper Combining | Shuffle 量削減 | Mapper 内で HashMap に保持し close 時に出力。 Combiner より確実 |
| Map-side Join | Shuffle 回避 | 小テーブルを DistributedCache に載せて Mapper で結合。 100 GB×1 GB なら有効 |
| Secondary Sort | Reducer 内処理高速化 | 複合キーで Reducer 入力をソート済にする。 集計順が決まる用途 |
| Speculative Execution | スラッカ対策 | 遅いタスクを別ノードで重複実行し早い方を採用。 既定で ON |
| Bloom Filter | 無駄な処理削減 | Join で右テーブルに存在しないキーを左で先に弾く。 Hive で自動適用 |
| 圧縮(Snappy/LZ4) | I/O 削減 | Map 出力、 最終出力を圧縮。 Snappy は CPU 軽め、 LZ4 は速度重視 |
| Parquet/ORC 列指向 | 読込量削減 | 必要列のみ読込み。 100 列中 5 列だけ使うなら 20 倍速い |
Hadoop は バッチ処理(一定期間溜めてから一括)に強いが、 リアルタイム性は弱い。 そこで派生したのが:
同じデータを 2 つの経路で処理:
複雑だが当時の現実解。 ヤフー、 Twitter、 LinkedIn が採用。
Kafka + Flink/Spark Streaming だけで完結。 ストリーミング処理が信頼性を持つようになった現代の主流。 Hadoop バッチ層を不要にする思想。
Hadoop は元々セキュリティが弱かったが、 2010 年代後半から大幅強化:
金融・医療・公共機関のオンプレ Hadoop はこれらを組合せて運用。 マイナンバー基盤など、 日本の機微情報処理にも使われている。
docker run apache/hadoop)。 HDFS の put/get/cat を試す10 問中 7 問以上 答えられれば実務レベル。 残りは関連用語ページを辿って補強。
| バージョン | リリース | 主要機能 |
|---|---|---|
| 0.20.x(初期) | 2009 | HDFS + MapReduce v1(JobTracker / TaskTracker)。 単一 NameNode で SPOF |
| 1.0.0 | 2011/12 | セキュリティ強化(Kerberos)、 HBase 連携強化 |
| 2.0.0 | 2012/05 | YARN 導入、 NameNode HA、 Federation、 Snapshot |
| 2.6.0 | 2014/11 | YARN ノードラベル、 透過的暗号化、 ローリングアップグレード |
| 3.0.0 | 2017/12 | イレイジャーコーディング(容量 50% 削減)、 Java 8、 YARN Timeline v2 |
| 3.2.0 | 2019/01 | S3A 改善、 GPU リソース管理、 Submarine(機械学習) |
| 3.3.0 | 2020/07 | ABFS 改善、 Java 11 対応、 多くのバグ修正 |
| 3.3.6 | 2023/06 | セキュリティ修正、 多くの安定化。 LTS 系列 |
| 3.4.0 | 2024/03 | Java 17 対応、 S3A 性能向上、 YARN 改善 |
SQL 風言語 HiveQL で MapReduce / Tez / Spark ジョブを生成。 メタストア(テーブル定義の DB)がエコシステムの中核に。 2008 年 Facebook 発祥。 当時 Facebook の SQL ユーザに Hadoop を使わせるためのブリッジとして。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 | -- Hive で SSDSE 風データを集計 CREATE EXTERNAL TABLE ssdse_b ( year INT, code STRING, prefecture STRING, population BIGINT, births BIGINT ) ROW FORMAT DELIMITED FIELDS TERMINATED BY ',' LOCATION '/user/data/ssdse-b/'; -- 都道府県別の出生数 TOP10 SELECT prefecture, SUM(births) AS total_births FROM ssdse_b WHERE year = 2023 GROUP BY prefecture ORDER BY total_births DESC LIMIT 10; |
HDFS の上に乗る列指向 NoSQL。 Google Bigtable のクローン。 ランダム読み書きが必要なケース(ユーザプロファイル、 リアルタイムダッシュボード等)に使う。 LSM tree 構造。
RDB(Oracle、 MySQL、 PostgreSQL)と HDFS 間のデータ移送ツール。 「sqoop import --connect ... --table users」で MySQL のテーブルが HDFS に。 ただし 2021 年に Apache Attic(メンテ終了)入り。 後継は Apache NiFi。
ログを Web サーバから HDFS へリアルタイム転送。 source → channel → sink の設定で柔軟に構成。 Web アクセスログ収集の定番だったが、 Kafka に置き換わるケースが増加。
LinkedIn 発祥の分散メッセージキュー。 Hadoop の「取込み」を担う標準ツール。 トピック単位で publish/subscribe、 ログを再生可能。 Lambda/Kappa アーキテクチャの基盤。
Hadoop ジョブのワークフロー管理。 XML で「MR ジョブ A 完了後に Hive ジョブ B、 失敗時に C」を記述。 近年は Airflow に押されて利用減。
分散システムの調整役。 NameNode HA、 HBase のリーダー選出、 Kafka のコンシューマグループ管理など、 Hadoop エコシステムのあちこちで使われる「裏方の主役」。
| データ | 行数 | サイズ | 適用技術 |
|---|---|---|---|
| SSDSE-B-2026(教育用) | 423 | 数百 KB | Excel、 pandas で十分 |
| SSDSE-A-2025(時系列) | 数千 | 数 MB | pandas、 R |
| e-Stat 国勢調査(マイクロデータ) | 数千万 | 数 GB | pandas + chunksize、 DuckDB |
| RESAS 地域経済データ(市町村×全産業) | 数億 | 数十 GB | Spark スタンドアロン、 Polars |
| 気象庁アメダス全国 10 年 | 数十億 | 数百 GB | Spark on YARN、 Hadoop 小クラスタ |
| Twitter / SNS ログ(数日) | 数百億〜兆 | 数 TB | Hadoop 大規模、 BigQuery |
| Web 全クロール(Common Crawl) | 数十億ページ | 数百 TB | Hadoop / Spark on AWS EMR |
| 天体望遠鏡データ(SKA) | 継続生成 | PB / 日 | 専用 HPC + Hadoop |
SSDSE-B-2026 は本物のビッグデータより 10 億倍小さい。 概念理解のために Hadoop を学ぶのは有益だが、 SSDSE 解析で実際に使う場面は無い。 「いつ必要になるか」を見極める眼が大事。
MapReduce の「Hello World」とされる WordCount。 入力テキストの中で各単語の出現回数を数える。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 | // Mapper public class TokenizerMapper extends Mapper<Object, Text, Text, IntWritable> { private final static IntWritable one = new IntWritable(1); private Text word = new Text(); public void map(Object key, Text value, Context context) { StringTokenizer itr = new StringTokenizer(value.toString()); while (itr.hasMoreTokens()) { word.set(itr.nextToken()); context.write(word, one); // (単語, 1) を出力 } } } // Reducer public class IntSumReducer extends Reducer<Text, IntWritable, Text, IntWritable> { private IntWritable result = new IntWritable(); public void reduce(Text key, Iterable<IntWritable> values, Context context) { int sum = 0; for (IntWritable val : values) sum += val.get(); result.set(sum); context.write(key, result); // (単語, 合計数) を出力 } } // Driver public class WordCount { public static void main(String[] args) throws Exception { Configuration conf = new Configuration(); Job job = Job.getInstance(conf, "word count"); job.setJarByClass(WordCount.class); job.setMapperClass(TokenizerMapper.class); job.setCombinerClass(IntSumReducer.class); // 同じ Reducer を Combiner に流用 job.setReducerClass(IntSumReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(IntWritable.class); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } } |
同じことを Hadoop Streaming(Python)で:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 | # Hadoop Streaming の mapper / reducer は、本来 sys.stdin から 1 行ずつ受け取る。 # ブラウザや対話環境には標準入力が無いので、ここでは行のリストを渡して同じ動きを再現する。 def mapper(lines): """1 行を単語に割り、(単語, 1) を出す。実際は print で標準出力に流す""" out = [] for line in lines: for word in line.strip().split(): out.append((word, 1)) return out def reducer(pairs): """単語で並べ替えたペアを受け取り、同じ単語の数を足し合わせる""" out = [] current, count = None, 0 for word, val in sorted(pairs): if word == current: count += int(val) else: if current is not None: out.append((current, count)) current, count = word, int(val) if current is not None: out.append((current, count)) return out lines = [ '北海道 東北 関東', '関東 中部 近畿', '関東 近畿 九州', ] pairs = mapper(lines) print('map の出力(先頭 5 件):', pairs[:5]) print('reduce の出力:', reducer(pairs)) |
同じことを Spark で書くと劇的に短い:
1 2 3 4 5 6 7 8 9 | from pyspark.sql import SparkSession spark = SparkSession.builder.appName('wc').getOrCreate() text = spark.read.text('input.txt') words = text.selectExpr('explode(split(value, " ")) AS word') counts = words.groupBy('word').count() counts.show() spark.stop() # 5 行で同じ処理。 これが Spark 普及の決定打 |
| パラメータ | 既定値 | 説明と推奨 |
|---|---|---|
| dfs.blocksize | 128 MB | HDFS ブロックサイズ。 大規模ファイル中心なら 256 MB ~ 1 GB に |
| dfs.replication | 3 | レプリカ数。 重要度の低いログは 2、 一時データは 1 も可 |
| dfs.namenode.handler.count | 10 | NameNode の RPC スレッド数。 大規模クラスタでは 100 ~ 200 |
| mapreduce.map.memory.mb | 1024 | Mapper の JVM メモリ。 OOM が出るなら 2048 ~ 4096 |
| mapreduce.reduce.memory.mb | 1024 | Reducer のメモリ。 集計が大きいなら増やす |
| mapreduce.job.reduces | 1 | Reducer 数。 出力データ量 / 1 GB 程度を目安に明示指定推奨 |
| mapreduce.task.io.sort.mb | 100 | Mapper の出力バッファ。 spill 回数削減のため 256 ~ 512 に |
| yarn.nodemanager.resource.memory-mb | 8192 | ノードあたりの YARN 利用可能メモリ。 物理メモリの 75% 程度に |
| yarn.scheduler.maximum-allocation-mb | 8192 | 1 コンテナ最大メモリ。 大きな Spark Executor のために 16 GB ~ 32 GB に |
| mapreduce.map.output.compress | false | Map 出力圧縮。 Shuffle ネックなら true(Snappy/LZ4)に |
MapReduce は関数型プログラミングの map と reduce を分散システムに持ち込んだものです。 数式で表すと次の 3 段階:
k₁ と値 v₁ を受け取り、 中間キー k₂ と中間値 v₂ のペアのリストを出す。 例:行 ID と都道府県データ → (地方コード, 人口) のリスト。 並列化が容易。k₂ を持つ v₂ を集約。 ここで ネットワーク経由のデータ転送 が大量に発生し、 ボトルネックになりやすい(Amdahl 法則の (1-p) 部分)。k₂ ごとに list(v₂) を最終出力 v₃ に変換。 例:地方ごとの合計人口を出す。SSDSE-B-2026 の例で言えば、 47 都道府県を 8 地方(北海道・東北・関東・中部・近畿・中国・四国・九州沖縄)に group by する処理は (都道府県コード, 人口) → (地方, 合計人口) の MapReduce そのものです。
シナリオ A:SSDSE-B-2026(47 都道府県 × 12 年 = 564 行、 約 200 KB) を Hadoop で集計するのは 明らかに過剰。 1 台の pandas で 0.1 秒で終わる規模。
シナリオ B:SSDSE-B を市区町村レベル × 全国小売店 POS × 10 年に拡張すると、 たとえば 1700 市町村 × 1000 店舗 × 5000 日 × 100 商品 = 約 850 億行、 1 行 200 byte と仮定して 17 TB。 ここで初めて Hadoop の出番。
HDFS ストレージ計算(架空のアクセスログ):1 日 100 GB のアクセスログを 1 年保管。 元データ $D = 100 \times 365 = 36.5$ TB。 レプリカ $R=3$ で実使用 109.5 TB。 ブロック数 $N = 36.5 \times 10^6 / 128 \approx 285{,}156$。 NameNode メモリ $\approx 285{,}156 \times 150$ byte ≒ 43 MB(余裕)。 ただしファイルを 1 KB の小ファイル 1 億個で持つと NameNode が破綻 — 小ファイル問題 の典型。
| データ規模 | 推奨ツール | 理由 |
|---|---|---|
| ~100 MB(SSDSE-B-2026 等) | pandas / Excel | 1 台で十分速い。 Hadoop 起動コストの方が大きい |
| 100 MB ~ 10 GB | pandas + chunksize / Polars / DuckDB | メモリ効率の良いシングルノード処理 |
| 10 GB ~ 1 TB | Spark(ローカル or 小クラスタ) | DataFrame API で慣れた書き方ができる |
| 1 TB ~ 1 PB | Hadoop / Spark on YARN / BigQuery / Snowflake | 分散処理が必須の領域 |
| 1 PB 超(ログ・センサー) | Hadoop + 専用ハードウェア / クラウド DWH | 本物の Hadoop ユースケース。 ヤフー、 メタ等の世界 |
合成データで 10 GB のファイルを 128 MB ブロックに分割した時の数とノード分散を計算する。
1 2 3 4 5 6 7 8 9 10 11 | import math file_mb = 10 * 1024 block_mb = 128 blocks = math.ceil(file_mb / block_mb) replicas = 3 total = blocks * replicas nodes = 20 per_node = total / nodes print(f"ブロック数: {blocks}") print(f"レプリカ込み: {total}") print(f"ノードあたり: {per_node} ブロック ({per_node * block_mb} MB)") |
💬 手計算 (Step 2) 12 ブロック/ノードと Python 出力が完全一致。
MapReduce の Map → Shuffle → Reduce を、 架空の短いテキストの ワードカウント(明記:下の単語列は説明用の 架空データ です)で 1 ステップずつ可視化します。 データを複数ノードに分割し、 各ノードが自分の近くのデータだけを数え、 同じ単語を同じ Reducer に集約して合計する — この 3 段を体感してください。
cat/dog/fish が同じ Reducer に集中し、 bird だけ別 Reducer になる。 一部の Reducer に仕事が偏ると全体が遅い Reducer に律速されます。HDFS は各ブロックを 別ノードに複製(既定 R=3、 下図は R=2)します。 ノードを停止してみて、 レプリカがある限りデータが読めることを確認してください。
① まず比較対象 — pandas でやる SSDSE-B-2026 集計(基準)
1 2 3 4 5 6 7 8 9 10 11 12 13 14 | import pandas as pd # SSDSE-B-2026 を 1 台の pandas で読み込み・集計(数百 KB なので一瞬) df = pd.read_csv('data/raw/SSDSE-B-2026.csv', skiprows=2, header=None, encoding='cp932') df = df.rename(columns={0: 'year', 2: 'prefecture', 3: 'population', 18: 'births'}) recent = df[df['year'] == 2023] # 都道府県別 平均出生数(実は 1 年 1 行なので合計と同じ) print('--- 上位 5 県(出生数) ---') print(recent.nlargest(5, 'births')[['prefecture', 'population', 'births']].to_string(index=False)) print('全国計(出生数):', recent['births'].sum()) print('行数:', len(recent), '/ メモリ:', df.memory_usage().sum() // 1024, 'KB') # → 数百 KB。 Hadoop は完全に過剰 |
② MapReduce の発想を pandas で再現(教育用ミニ実装)
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 | import pandas as pd from collections import defaultdict df = pd.read_csv('data/raw/SSDSE-B-2026.csv', skiprows=2, header=None, encoding='cp932') df.columns = ['year', 'code', 'prefecture'] + [f'col{i}' for i in range(3, df.shape[1])] recent = df[df['year'] == 2023][['prefecture', 'col18']].rename(columns={'col18': 'births'}) # Map フェーズ:行を (キー, 値) に変換 def mapper(row): pref = row['prefecture'] # 地域分類(東日本 / 西日本)をキーに east = ['北海道', '青森県', '岩手県', '宮城県', '秋田県', '山形県', '福島県', '茨城県', '栃木県', '群馬県', '埼玉県', '千葉県', '東京都', '神奈川県', '新潟県', '富山県', '石川県', '福井県', '山梨県', '長野県', '岐阜県', '静岡県', '愛知県'] return ('東日本' if pref in east else '西日本', row['births']) mapped = [mapper(row) for _, row in recent.iterrows()] print('Map 出力(最初5件):', mapped[:5]) # Shuffle フェーズ:キーでグループ化 shuffled = defaultdict(list) for k, v in mapped: shuffled[k].append(v) # Reduce フェーズ:キーごとに集約 reduced = {k: sum(v) for k, v in shuffled.items()} print('Reduce 結果:', reduced) # → これが MapReduce の最小モデル。 ノードに分散するなら mapped を分割して並列実行 |
③ PySpark で同じ集計(Hadoop 互換 API)
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 | from pyspark.sql import SparkSession from pyspark.sql import functions as F # ローカルモードで起動(クラスタ無しでも動く) spark = SparkSession.builder.appName('ssdse_demo').master('local[*]').getOrCreate() # HDFS にあるなら 'hdfs:///data/SSDSE-B-2026.csv'。 ローカルなら通常パス df = spark.read.csv('data/raw/SSDSE-B-2026.csv', header=False, inferSchema=True, encoding='cp932') df = df.toDF(*([f'col{i}' for i in range(df.columns.__len__())])) # 年度フィルタ・都道府県別集計 recent = df.filter(F.col('col0') == 2023) agg = recent.groupBy('col2').agg(F.sum('col18').alias('total_births')) agg.orderBy(F.desc('total_births')).show(5) # 実行計画(Catalyst optimizer)を見る agg.explain() spark.stop() |
④ HDFS を CLI で操作(参考。 Hadoop 環境が必要)
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 | # 注意: hadoop コマンドは Hadoop 環境(CDH/HDP/EMR 等)でのみ動作 # HDFS にローカルファイルをアップロード # $ hadoop fs -put data/raw/SSDSE-B-2026.csv /user/me/ # 中身を確認 # $ hadoop fs -cat /user/me/SSDSE-B-2026.csv | head # ブロック配置(どの DataNode に複製されたか)を確認 # $ hdfs fsck /user/me/SSDSE-B-2026.csv -files -blocks -locations # レプリカ数を変更(既定 3 → 2 に減らしてストレージ節約) # $ hadoop fs -setrep -w 2 /user/me/SSDSE-B-2026.csv # 容量確認 # $ hadoop fs -df -h / # $ hadoop fs -du -h /user/me/ |
⑤ Amdahl の法則をシミュレーション
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 | import numpy as np import pandas as pd def amdahl_speedup(f, p): """並列化可能割合 f, ノード数 p のときの理論速度向上""" return 1.0 / ((1 - f) + f / p) rows = [] for f in [0.5, 0.8, 0.9, 0.95, 0.99]: for p in [1, 2, 4, 8, 16, 32, 64, 128, 256, 512, 1024]: rows.append({'f': f, 'p': p, 'speedup': amdahl_speedup(f, p)}) result = pd.DataFrame(rows).pivot(index='p', columns='f', values='speedup').round(2) print('--- 速度向上倍率(Amdahl)---') print(result) # 教訓: f=0.99 でも p=1024 で 91 倍止まり。 100 倍が理論上限 |
🎯 このコードでやること:SSDSE-B-2026 の 47 都道府県データを「分散処理風」に分割し、 並列化可能割合 p=0.9 と p=0.7 の 2 通りで Amdahl 法則のスピードアップ曲線を計算する。
📥 入力データ:SSDSE-B-2026(47 都道府県 × 約 80 変数)。 各都道府県を 1 ブロック相当として扱う想定。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 | import numpy as np import pandas as pd # Amdahl 法則の定義 def amdahl_speedup(N, p): """並列度 N、 並列化可能割合 p のスピードアップ S(N)""" return 1.0 / ((1.0 - p) + p / N) # SSDSE 47 県を分散ブロックと見立てる df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=[1]) n_blocks = len(df) # 47 都道府県 # 2 つのシナリオ N_list = [1, 2, 4, 8, 16, 32, 47, 100, 1000] for p in [0.9, 0.7]: print(f"--- 並列化可能割合 p = {p} ---") for N in N_list: s = amdahl_speedup(N, p) print(f" N={N:>4d} ノード → S(N)={s:.2f}x") print(f" 上限 S(∞) = {1/(1-p):.2f}x\n") </div> |
📤 実行結果:
💬 結果の読み方:p=0.9 でも 47 ノードで 8.3 倍止まり。 p=0.7 だと 3.2 倍が事実上の天井。 SSDSE のような数十 MB 規模で Hadoop を使う場合、 逐次部分(シャッフル・ジョブ起動)が支配的 になり、 N を増やしても無駄が大きい。 これが「中小データに Hadoop は不向き」の数値的根拠です。
🎯 このコードでやること:SSDSE-B-2026 の 47 都道府県データを PySpark で読み込み、 「地方ごとの合計人口」を MapReduce 風に集計する。 ローカルモード(master='local[*]')で動作確認。
📥 入力データ:SSDSE-B-2026 から「都道府県コード」「都道府県名」「総人口」の 3 カラム。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 | from pyspark.sql import SparkSession from pyspark.sql.functions import sum as F_sum, col, when # Spark セッション起動(ローカルモード、 全 CPU 使用) spark = SparkSession.builder \ .appName('SSDSE-Hadoop-Demo') \ .master('local[*]') \ .getOrCreate() # SSDSE-B-2026 を Spark DataFrame で読込 sdf = spark.read.csv( 'data/raw/SSDSE-B-2026.csv', header=True, encoding='cp932', inferSchema=True ) # 2 行目(説明行)を除去 header_row = sdf.first() sdf = sdf.filter(col('SSDSE-2026') != header_row['SSDSE-2026']) # 都道府県コード先頭 2 桁から地方を分類(Map ステップに相当) sdf = sdf.withColumn('region', when(col('SSDSE-2026').substr(2, 2) <= '01', '北海道') .when(col('SSDSE-2026').substr(2, 2) <= '07', '東北') .when(col('SSDSE-2026').substr(2, 2) <= '14', '関東') .when(col('SSDSE-2026').substr(2, 2) <= '23', '中部') .when(col('SSDSE-2026').substr(2, 2) <= '30', '近畿') .when(col('SSDSE-2026').substr(2, 2) <= '35', '中国') .when(col('SSDSE-2026').substr(2, 2) <= '39', '四国') .otherwise('九州沖縄') ) # 地方ごとに合計人口を集計(Reduce ステップ) result = sdf.groupBy('region').agg(F_sum('総人口').alias('合計人口')) \ .orderBy(col('合計人口').desc()) result.show() </div> |
📤 実行結果:
💬 結果の読み方:関東 4,380 万人で日本の人口集中が顕著。 PySpark で書いても処理量は同じだが、 数 TB 規模になればクラスタが効く。 SSDSE 規模(数十 MB)では pandas の groupby が桁違いに速い。
🎯 このコードでやること:SSDSE-B-2026 を dask.dataframe で読み込み、 同じ集計を「ノード数 4」想定で実行。 pandas API と互換のため移行コストが低いことを確認。
📥 入力データ:SSDSE-B-2026 の同データを dask.dataframe で 4 パーティションに分割。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 | import dask.dataframe as dd import pandas as pd import time # SSDSE-B-2026 を pandas で読込 df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=[1]) # dask DataFrame に変換(4 パーティション) ddf = dd.from_pandas(df, npartitions=4) # 地方分類関数 def to_region(code): code_int = int(str(code)[1:3]) if code_int == 1: return '北海道' elif code_int <= 7: return '東北' elif code_int <= 14: return '関東' elif code_int <= 23: return '中部' elif code_int <= 30: return '近畿' elif code_int <= 35: return '中国' elif code_int <= 39: return '四国' else: return '九州沖縄' ddf['region'] = ddf['SSDSE-2026'].map(to_region, meta=('region', 'object')) # 地方ごとに集計(map → reduce 相当) t0 = time.time() result = ddf.groupby('region')['A1101'].sum().compute() print(f'Elapsed: {time.time()-t0:.3f} sec') print(result.sort_values(ascending=False)) </div> |
📤 実行結果:
💬 結果の読み方:dask は pandas と同じ API。 47 行では並列化の効果はほぼなく、 むしろ Spark 起動のオーバヘッドが目立つ(pandas の groupby: 約 0.003 秒)。 並列化の恩恵が出るのは 数億行以上。 「規模に応じてツールを選ぶ」のがエンジニアリングの基本。
HDFS の特徴は レプリケーション(既定 3 系統)です。 1 つのブロックを別ノードに 3 つコピーしておくことで、 1 ノードが落ちてもデータが失われない。 ノード故障確率を q、 レプリカ数を r とすると、 ブロックが完全に失われる確率は:
$$ P_{loss} = q^r $$q=0.02、 r=3 のとき P_loss = 0.02³ = 0.000008(百万分の 8)。 100 万ブロックでも期待損失 8 ブロック。 これが「HDFS は安心して TB 級データを置ける」根拠です。 一方で r=1 にすると q=0.02 そのままなので、 1 万ブロック中 200 ブロックが消失する計算。
| レプリカ数 r | P_loss(q=0.02) | 100 万ブロック中の期待損失 |
|---|---|---|
| 1 | 0.02000 | 20,000 ブロック |
| 2 | 0.00040 | 400 ブロック |
| 3(既定) | 0.000008 | 8 ブロック |
| 4 | 0.00000016 | 0.16 ブロック |
| 年 | 出来事 | 技術的意義 |
|---|---|---|
| 2003 | Google GFS 論文 | 分散ファイルシステムの設計指針が公開 |
| 2004 | Google MapReduce 論文 | 関数型 + 分散コンピューティングの提示 |
| 2006 | Hadoop プロジェクト分離 | Doug Cutting が Yahoo! でフルタイム開発 |
| 2009 | Cloudera 創業 | エンタープライズ商用化の幕開け |
| 2011 | Hadoop 1.0 リリース | 本格的な production 採用が拡大 |
| 2013 | YARN(Hadoop 2.0) | リソース管理分離、 MapReduce 以外のフレームも乗る |
| 2014 | Spark 急上昇 | インメモリ処理で 10〜100 倍高速、 Hadoop の主役交代 |
| 2017 | クラウド移行加速 | S3 + EMR、 GCS + Dataproc が主流に |
| 2020 | Cloudera 上場廃止 | Hadoop 商用市場の収縮を象徴 |
| 2022 | Iceberg / Delta Lake 普及 | HDFS を介さない data lake が標準 |
| 2026 | 「Hadoop 思想」の継承 | Spark + クラウド + Iceberg が主流、 Hadoop は教養化 |
JOIN や ORDER BY を書くと、 裏で巨大な Shuffle が発生。 SORT BY(パーティション内ソート)と ORDER BY(全件ソート)の違いを知らずに後者を使い、 Reducer 1 個でクラスタ停止級に詰まる事故が多い。1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 | # ファイル一覧 hadoop fs -ls /user/me/ hadoop fs -ls -R /user/me/ # 再帰 # ローカル → HDFS hadoop fs -put local.csv /user/me/ hadoop fs -copyFromLocal local.csv /user/me/ # HDFS → ローカル hadoop fs -get /user/me/data.csv ./ hadoop fs -copyToLocal /user/me/data.csv ./ # 中身確認 hadoop fs -cat /user/me/data.csv | head hadoop fs -tail /user/me/log.txt # 削除 hadoop fs -rm /user/me/old.csv hadoop fs -rm -r /user/me/old_dir/ # 再帰 hadoop fs -rm -skipTrash /user/me/x # ゴミ箱経由せず即削除 # ディレクトリ作成・名前変更 hadoop fs -mkdir /user/me/newdir hadoop fs -mv /user/me/old.csv /user/me/new.csv # 容量・ブロック情報 hadoop fs -df -h / hadoop fs -du -h /user/me/ hdfs fsck /user/me/data.csv -files -blocks -locations # レプリカ数変更 hadoop fs -setrep -w 2 /user/me/data.csv # パーミッション hadoop fs -chmod 755 /user/me/data.csv hadoop fs -chown me:mygroup /user/me/data.csv |
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 | # 実行中ジョブ一覧 yarn application -list # ジョブ詳細 yarn application -status application_1234_0001 # ジョブ強制終了 yarn application -kill application_1234_0001 # ジョブログ取得 yarn logs -applicationId application_1234_0001 # キュー情報 yarn queue -status default # ノード情報 yarn node -list -all |
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 | # NameNode セーフモード(メンテ時) hdfs dfsadmin -safemode get hdfs dfsadmin -safemode enter hdfs dfsadmin -safemode leave # データの整合性チェック hdfs fsck / -files -blocks hdfs fsck / -delete # 破損ブロックを削除 # データバランサ起動(DataNode 間のリバランス) hdfs balancer -threshold 10 # NameNode メタデータバックアップ hdfs dfsadmin -saveNamespace hdfs dfsadmin -fetchImage /backup/ # クラスタ全体の HA 状態 hdfs haadmin -getServiceState nn1 hdfs haadmin -getServiceState nn2 hdfs haadmin -failover nn1 nn2 # 手動フェイルオーバ |
| 方式 | 月額コスト目安 | 前提・特徴 |
|---|---|---|
| オンプレ Hadoop(自社運用) | 100~300 万円 | 100 ノード(各 10 TB)、 レプリカ 3 で実 333 TB。 電気代・運用人件費別 |
| AWS S3 + EMR | 300~600 万円 | S3 標準 $0.023/GB/月 × 1 PB ≒ $23,000/月 + EMR コンピュート |
| GCP Cloud Storage + Dataproc | 300~500 万円 | GCS Standard $0.020/GB/月。 Dataproc 起動が速い |
| BigQuery | 200~400 万円 | ストレージ $0.020/GB/月、 クエリ $5/TB(処理量課金) |
| Snowflake | 300~700 万円 | ストレージ $40/TB/月 + コンピュート(クレジット課金) |
近年の傾向:オンプレ Hadoop は 運用人件費(年数千万)と更新(3 年で機材入替)を入れるとクラウドより高くなるケースが大半。 特に Cold データ(年 1 回しか読まない)は S3 Glacier 等の低価格層に置く方が圧倒的に安い($0.004/GB/月)。
Apache Mahout が標準でした。 K-means、 ALS、 ナイーブベイズ、 ランダムフォレストなどを MapReduce で実装。 ただし反復処理に弱く、 100 回ループする勾配降下は実用に耐えませんでした。
機械学習は Spark MLlib、 scikit-learn、 PyTorch/TensorFlow が主流。 Hadoop は 「ML のためのデータ前処理・保管庫」として残る。 ETL(抽出・変換・ロード)の中核として、 Spark on Hadoop でデータを整え、 Parquet で保存して、 ML フレームワークから読み出す、 という構成が多い。
Petastorm(Uber 開発)を使うと Parquet データを PyTorch / TensorFlow から効率的に読める。 これにより HDFS 上の大規模学習データを直接 GPU クラスタに供給可能。
1 2 3 4 5 6 7 8 9 | # PyTorch から Hadoop 上の Parquet を読む例 from petastorm.pytorch import DataLoader from petastorm import make_batch_reader with make_batch_reader('hdfs:///data/training.parquet') as reader: train_loader = DataLoader(reader, batch_size=64) for batch in train_loader: # batch は PyTorch tensor pass |
結論:Hadoop の「思想」(分散・レプリケーション・データローカリティ)は普遍的。 ただし具体的な実装としては Spark + クラウドストレージ を主軸に学ぶのが効率的。 Hadoop はそのバックボーンとして「概念で押さえる」のが 2026 年のスタンスです。
┌─ ストレージ層:HDFS(中心)/S3/GCS/Azure Blob
├─ リソース管理:YARN/Kubernetes/Mesos
├─ 計算エンジン:MapReduce(旧)/Spark/Flink/Tez
├─ SQL レイヤ:Hive/Presto/Trino/Impala
├─ NoSQL:HBase/Cassandra
├─ 取込み:Kafka/Flume/Sqoop/NiFi
├─ ワークフロー:Airflow/Oozie
├─ テーブルフォーマット:Parquet/ORC/Avro/Iceberg/Delta Lake/Hudi
└─ セキュリティ:Kerberos/Ranger/Knox/Sentry
「Hadoop」は単独で完結する手法ではなく、 隣接領域と連携することで真価を発揮する。
上流の分散ファイルシステム (HDFS) でデータ配置を設計し、 並列の Spark・Flink とインメモリ性能を比較し、 下流の Hive・Pig・Presto でクエリインタフェースを提供する。 Hadoop は MapReduce 時代の標準だが、 現代では Spark への移行・クラウドネイティブ (S3 + EMR) への置き換えが進む過渡期技術として理解する。
「Hadoop」を実際の課題に当てはめるとき、 状況別に何を選ぶかを 3 段階で判定する。
Hadoop は 2010 年代の分散処理標準だが、 2025 年現在は Spark とクラウドネイティブが主流。 新規プロジェクトで Hadoop を選択する場面は少なく、 既存資産の保守・段階的移行のための知識として位置づける。
本ページ上部の WordCount 体験(🎮 触って理解する)は「単語→出現回数」の集約でした。 ここではもう一歩踏み込み、 実在データ SSDSE-B-2026 の 47 都道府県を「8 地方」に畳み込む集計 を題材に、 MapReduce の Map→Shuffle→Reduce がなぜ WordCount と同じ形なのか、 そしてそこに潜む実務の罠を掘り下げます。
WordCount では 単語 がキーでした。 地方別人口集計では 地方名 がキーになります。 構造は完全に同型です:
(地方名, 総人口) を 1 件ずつ emit する。 47 行がそれぞれ別ノードで並列処理できる(=並列化可能部分 $p$)。実際に data/raw/SSDSE-B-2026.csv(2023 年・47 都道府県、 列 A1101=総人口)で Reduce の出力を計算すると、 次のようになります(実測値):
| Reduce キー(地方) | 担当県数(Map 入力数) | 総人口 合計〔人〕 | 全国比 |
|---|---|---|---|
| 北海道 | 1 | 5,092,000 | 4.1% |
| 東北 | 6 | 8,318,000 | 6.7% |
| 関東 | 7 | 43,527,000 | 35.0% |
| 中部 | 9 | 20,749,000 | 16.7% |
| 近畿 | 7 | 21,990,000 | 17.7% |
| 中国 | 5 | 7,070,000 | 5.7% |
| 四国 | 4 | 3,578,000 | 2.9% |
| 九州沖縄 | 8 | 14,029,000 | 11.3% |
| 全国計 | 47 | 124,353,000 | 100% |
※ 数値は pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=[1]) で 2023 年行を抽出し、 A1101(総人口)を地方別に groupby().sum() した実測値。 地方区分は総務省の 8 地方区分に準拠(東京都単独の総人口は 14,086,000 人、 最小の鳥取県は 537,000 人)。
上の表をよく見てください。 関東キー 1 つで全国の 35.0%(43,527,000 人分) を占めます。 一方で四国キーは 2.9%(3,578,000 人分)に過ぎません。 MapReduce では 1 つの Reduce キーは 1 つの Reducer に必ず割り当たるため、 キー間の偏りがそのまま Reducer 間の負荷偏りになります。
the のような超高頻度語が 1 Reducer に集中するのと、 構造的に完全に同じ現象です。 キーの分布が偏るデータでは、 ノードを増やしても頭打ちどころか、 偏ったキーの担当ノードが単独のボトルネックになります。要点:「並列度=速度」ではない。 キー分布が偏った瞬間、 台数を増やしても最遅 Reducer が壁になる。 スケールアウトの前に「集約キーの分布」を必ず確認せよ — これは Hadoop でも Spark でも BigQuery でも変わらない普遍則です。
Combiner が使えるのは演算が結合律・交換律を満たすときだけです。 総人口の sum は満たすので、 各 Map ノードで先に県内合計を作り「地方名→部分和」を emit すれば、 Shuffle で流れるレコードは激減します(47 件 → 各ノードの部分和のみ)。 一方 中央値やユニーク数(distinct count)は結合律を満たさないため Combiner をそのまま適用できず、 スキュー対策が一段難しくなります。 「その集約が Combiner 可能か?」は分散集計を設計する最初のチェックポイントです。
Spark ではこの Map→Shuffle→Reduce が reduceByKey / groupBy().agg() として同じ骨格で現れ、 reduceByKey は Combiner 相当の Map 側事前集約を自動で行います(対して groupByKey は全件を Shuffle するため非推奨)。 つまり本ページで学んだ「偏りキー=ストラグラー」という直感は、 実行エンジンが Hadoop から Spark に替わってもそのまま通用します。 なお SSDSE-B のような数 MB の CSV では、 そもそも 1 台の pandas groupby が瞬時に終わるため分散処理は過剰であり、 上の表も pandas で 1 秒未満で計算できます — 分散が意味を持つのは同型の集計を TB~PB 級で回すときだけ、 という判断軸(📍 あなたが今見ているもの)も併せて思い出してください。
reduceByKey のスキュー対策が本ページの続き。