「分散データ処理」は単一マシンで処理しきれないデータ量・計算量を、 複数ノードで並列処理する技術群。 MapReduce → Spark → Dask → Ray と進化し、 ストリーミング (Kafka/Flink) も中核に含む。 本ページでは MapReduce のパラダイム・Spark の RDD/DataFrame・Dask のタスクグラフ・データシャーディング・耐障害性 (lineage)・shuffle のコストを整理する。
これらのキーワードは「データを分割 → ノードで並列処理 → shuffle で結果統合 → 障害から回復」という分散処理の中核 4 ステップを構成する。
🍰 まずはやさしく
みんなで分担して計算する方法です。
大量のデータを早く処理するために使います。
スマホのアプリで膨大な情報を扱うときに役立ちます。
この章では分散処理の仕組みと代表的な技術を読みます。
データの分散処理 = データを複数のマシン (ノード) に分割して並列に計算させる方式。1 台では時間・メモリが足りないデータを、たくさんの普通のマシンで分担する。
🍰 まずはやさしく
AIなどの大きなシステムで使う考え方です。
1台のパソコンでは足りない計算をさせるために使います。
47都道府県のデータを分担して集計する例で考えます。
この章ではデータの置き場所と速度の関係について読みます。
用語集 → ビッグデータ / AI インフラ 分野 → データの分散処理 (Distributed Data Processing)。AI クラウド 教材で頻出する横串の概念です。深層学習で 1 GPU では足りないとき、ペタバイト級のログを 1 台では処理できないときに登場します。
本ページの SSDSE 文脈では 47 都道府県 = 47 ワーカのメタファーを使って、Pythonの concurrent.futures / multiprocessing / dask で並列処理する例を見せます。
分散処理の鉄則「計算をデータの近くへ送る」を Spark / Hadoop は 5 段階で評価する。 47 都道府県シャードを各リージョンに置く設計だと、 PROCESS_LOCAL でない場合に遅延が桁違いに伸びる。
| 局所性レベル | 説明 | 読込速度 |
|---|---|---|
| PROCESS_LOCAL | 同一 JVM プロセス内 (キャッシュ済み) | 数 GB/s |
| NODE_LOCAL | 同一マシン内 (別 JVM またはディスク) | 数百 MB/s (SSD) |
| RACK_LOCAL | 同一ラック内 (LAN) | 数十 MB/s |
| ANY (DataCenter) | 同一 DC、 別ラック (LAN) | 数 MB/s |
| RemoteDC | 別データセンタ (WAN) | 数百 KB/s |
SSDSE 設計例: 47 都道府県データを「東京リージョン (本州 27 県)」「大阪リージョン (西日本 14 県)」「札幌リージョン (北海道+東北 6 県)」に分散したら、 「全国集計」は RemoteDC 通信が発生し WAN 帯域が律速になる。 解決策は「事前に各リージョンで部分集約 (combiner) し、 中央に小さな集計値だけ送る」。
「ノード故障」は分散システムでは毎日起こる前提。 1000 ノードで MTBF (平均故障間隔) 3 年のサーバを使うと、 毎日 1 台は壊れる計算。 47 都道府県ノードでも、 月 1 台程度の故障は想定すべき。
| 失敗モード | 原因例 | 対策 |
|---|---|---|
| クラッシュ故障 | OS パニック、 電源断 | レプリカ + リーダー再選出 |
| オミッション故障 | メッセージ消失 | タイムアウト + リトライ |
| タイミング故障 | 遅延が予想外に大きい | SLA 監視 + 速い完了率の冗長化 (hedge request) |
| レスポンス故障 | 間違った値を返す | チェックサム、 投票 |
| Byzantine 故障 | 悪意 / バグで矛盾した値 | PBFT、 BLS 署名 |
| ネット分断 (Partition) | スイッチ障害、 ケーブル切断 | CAP に従い CP or AP を選択 |
| スプリットブレイン | 分断で複数リーダー | クォーラム、 STONITH (フェンシング) |
| カスケード故障 | 1 ノード故障 → 隣接過負荷 → 連鎖 | サーキットブレーカ、 backpressure |
| Thundering Herd | 復旧時に全クライアントが同時アクセス | jittered exponential backoff |
| スローノード (Straggler) | 1 ノードだけ遅い | 投機的実行 (speculative execution) |
| 時刻ずれ (Clock Skew) | NTP 同期失敗 | 論理時計 (Lamport, Vector Clock) |
| グレイ故障 | 部分的に動作 (検知困難) | end-to-end ヘルスチェック、 確率的フェイルアウト |
| 書籍 / 論文 | 著者 | レベル | 推薦点 |
|---|---|---|---|
| Designing Data-Intensive Applications | M. Kleppmann | 中級 | 分散システムの教科書 (略称 DDIA) |
| Distributed Systems | A. Tanenbaum & M. van Steen | 中級 | 理論面の定番 |
| MapReduce: Simplified Data Processing... (2004) | J. Dean & S. Ghemawat (Google) | 中級 | 分散処理の起源論文 |
| The Google File System (2003) | S. Ghemawat et al. | 中級 | HDFS の元論文 |
| Spanner: Google's Globally-Distributed DB (2012) | J. Corbett et al. | 上級 | TrueTime と外部整合性 |
| In Search of an Understandable Consensus Algorithm (Raft) | Ongaro & Ousterhout (2014) | 中級 | Paxos を理解しやすく |
| Dynamo: Amazon's Highly Available Key-value Store | G. DeCandia et al. (2007) | 上級 | consistent hashing + vector clock |
| Bigtable: A Distributed Storage System... | F. Chang et al. (2008) | 中級 | HBase / Cassandra のルーツ |
| Resilient Distributed Datasets (RDD) | M. Zaharia et al. (2012) | 上級 | Spark の起源論文 |
| Site Reliability Engineering | Google SRE Team | 中級 | 分散システムの運用論 |
| Hadoop: The Definitive Guide | T. White | 初級 | Hadoop エコシステム入門 |
| Learning Spark, 2nd ed. | J. Damji et al. | 初級 | Spark 入門書 |
分散処理の効果は「データ量 × 計算量」が大きいほど顕著になる。 SSDSE-B-2026 の 47 行 × 約 100 列のデータを 都道府県別にレプリケートして 47,000 行に水増しし、 単一 pandas と Dask 並列 (4 worker) で集計時間を比較する。 これは strong scaling (同一問題の高速化) のミニ実験である。
このコードでやること: SSDSE-B-2026 を読み込み、 1000 倍に複製した上で、 都道府県別の人口・出生数の集計を pandas と dask.dataframe で実施し、 実行時間を比較する。
📥 入力データ (SSDSE-B-2026 抜粋):
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 | import pandas as pd import dask.dataframe as dd import time # SSDSE-B-2026 読み込み (cp932) → 最新 2023 年の 47 都道府県だけ抽出 df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=[1]) df = df[df['SSDSE-B-2026'] == 2023] # 47 都道府県 × 1 年 = 47 行 print(f'元データ: {len(df)} 行 × {len(df.columns)} 列') # 1000 倍に複製して 47,000 行のミニビッグデータ化 df_big = pd.concat([df] * 1000, ignore_index=True) print(f'拡張後: {len(df_big):,} 行') # (1) pandas で集計 t0 = time.time() agg_pd = df_big.groupby('Prefecture').agg({'A1101': 'mean', 'A4101': 'sum'}) t_pd = time.time() - t0 # (2) Dask で集計 (4 partition に分割 = 並列度 4) ddf = dd.from_pandas(df_big, npartitions=4) t0 = time.time() agg_dk = ddf.groupby('Prefecture').agg({'A1101': 'mean', 'A4101': 'sum'}).compute() t_dk = time.time() - t0 print(f'pandas: {t_pd*1000:.1f} ms') print(f'Dask×4: {t_dk*1000:.1f} ms') print(f'速度比: {t_pd/t_dk:.2f}x') |
📤 実行例(1 台の PC での一例。値は環境と実行ごとに変わる):
💬 速くなりません。むしろ 20 倍遅くなります。 47,000 行の groupby は pandas なら 1 ms 前後で終わってしまい、 Dask 側はタスクグラフの構築・分割・結果の結合だけで 20 ms 前後を使うためです。 Amdahl の法則 $S \le 1 / ((1-p) + p/N)$ は「並列化できる部分 $p$ の割合」しか見ておらず、 分散処理そのものの固定コストを含みません。 このコストはデータ量にほとんど依存しないので、 データが小さいほど不利になります。 分散処理が効いてくるのは、 1 台のメモリに載らない規模(数十 GB 以上)や、 集計 1 回に数秒以上かかる規模からです。 「分散にすれば速くなる」わけではない——これがこの実測から読み取るべきことです。 なお Dask の初回実行はウォームアップを含むため 100 ms 超になることがあり、 2 回目以降が 20 ms 前後に落ち着きます。
分散処理の評価指標として、 strong scaling (問題サイズ固定で worker 数を増やしたときの高速化) と、 weak scaling (worker 数比例でデータ量も増やしたときの実行時間維持率) の 2 つを使い分ける。 strong scaling は Amdahl の法則、 weak scaling は Gustafson の法則で議論する。
Gustafson の法則 (数式を言葉で読み解く): $S(N) = N - \alpha(N-1)$ ここで $\alpha$ は逐次実行部分の割合。 worker 数 $N$ が増えても、 逐次部分の重みが下がり続けるため、 strong scaling より楽観的な見積もりとなる。 例えば $\alpha$=0.05, $N$=100 なら $S$=100 − 0.05×99 = 95.05 倍。 ビッグデータ解析では worker 投入と同時に対象データ量も拡大できるため、 weak scaling 視点が現実的である。
| 指標 | 定義 | 理想値 | 実務典型値 |
|---|---|---|---|
| strong scaling 効率 $E_s$ | $T_1 / (N \cdot T_N)$ | 1.0 (= 100%) | 0.6 - 0.8 |
| weak scaling 効率 $E_w$ | $T_1 / T_N$ (データ比例) | 1.0 | 0.8 - 0.95 |
| スループット (rec/s) | 単位時間処理レコード数 | 線形増加 | 準線形 |
| shuffle データ量 | worker 間転送バイト | 0 (= map only) | 入力の 0.5 - 2x |
groupby, join, sort は worker 間で大量データを再配置する必要があり、 ネットワーク帯域がボトルネックになる。 partition 設計でキーを揃えると軽減できる。ここからは SSDSE-B-2026 (47 都道府県 × 約 100 列の公的統計データ) を使い、 「単一マシンで完結できるサイズ」と「分散処理が必要になるサイズ」の境目を**実測値で**追っていく。 分散処理は万能の高速化策ではない。 むしろ小規模データでは遅くなる。 どこからが「分散の出番」かを、 工程ごとに具体的な数字と図表で示すのがこのセクションの目的である。
分散処理は「スケジューリング」「タスク分割」「データ転送 (shuffle)」「結果集約」のオーバーヘッドを必ず伴う。 SSDSE-B-2026 のように 47 行しかないデータでは、 これらのオーバーヘッドが計算本体 (たかだか数 µs) を桁違いに上回る。 下表は、 単一マシン pandas, Dask (local cluster), PySpark (local mode) で SSDSE-B-2026 の「総人口の県別合計」を実行したときの大小関係を示す説明用の仮の数値で、 実測値ではない(この教材の中では Dask と PySpark を同梱していないため再現もできない)。 読み取ってほしいのは個々の数字ではなく、 「小さなデータでは分散処理のほうが 桁違いに遅い」という関係のほうである。 お手元の環境で測れば、 数字は違っても同じ傾向が出る。
| 実行系 | 中央値 (ms、 仮の値) | 分散数 | 起動オーバーヘッド (ms) | 備考 |
|---|---|---|---|---|
pandas (df.groupby().sum()) | 0.18 | 1 プロセス | 0 | 基準 |
| Dask local (4 worker) | 14.6 | 4 partition | 12 (scheduler) | 約 80 倍遅い |
| PySpark local (4 thread) | 92.1 | 4 partition | 75 (JVM warm) | 約 500 倍遅い |
| Dask local (1 worker) | 8.4 | 1 partition | 8 | 並列化しても効果なし |
この仮の表は分散処理の本質を示している。 並列化の利益 (約 4 倍) より、 タスクをまたぐ通信・スケジューリングのコスト (12-75 ms) の方が圧倒的に大きいのだ。 SSDSE-B-2026 のようなマイクロデータでは、 pandas の C 実装 (numpy のベクトル演算 + Cython の groupby) が常勝する。 では、 どこから分散処理が pandas を上回るのか? 次の図で、 実際に測って確かめる。
SSDSE-B-2026 (2023 年度の 47 行、 都道府県名と総人口・65 歳以上人口・出生数・死亡数・消費支出の 5 列) を 1〜300,000 倍に複製し (= 47 行 〜 1,410 万行)、 同じ groupby('Prefecture').sum() を pandas と Dask (dask.dataframe、 scheduler='threads'、 8 partition) で 5 回ずつ測った中央値 (横軸: 行数 log, 縦軸: 時間 ms log)。 処理時間は実行する環境 (CPU・メモリ・ライブラリの版) で大きく変わるので、 数値は 1 台の PC での目安である。 この測定では、 1,410 万行まで増やしても Dask が pandas より速くなる交差点は現れなかった。 Dask / pandas の時間の比は 47 行で約 66 倍、 47 万行で 1.84 倍、 1,410 万行で 1.09 倍と、 行数が増えるほど縮まっていく。
図 1 から読み取れる教訓は 3 つある。 (1) 小規模では pandas が圧勝: 47 行〜4,700 行では、 pandas の処理時間 (0.1〜0.15 ms) は Dask がタスクを組んで配るだけの固定費 (約 6〜7 ms) にまったく届かない。 (2) 行数が増えると差は縮まる: 47 万行を超えると固定費は処理本体に埋もれ、 1,410 万行では Dask の方が 9% 遅いだけになる。 (3) メモリに収まる量なら 1 台の pandas で十分: この測定ではメモリに収まる範囲で交差点が現れなかった。 Dask が本当に効くのは、 1 台のメモリに載らない量 (pandas がメモリ不足で動かなくなる量) を partition ごとに読み込んで処理するときや、 複数台に計算を分けるときである。 1 億行規模の測定はこのページでは行っていない。
分散処理の対象は、 多くの場合「大量のレコードに対する集計・変換・結合」である。 では SSDSE-B-2026 で集計対象となる「総人口」はどのような分布なのか? 47 都道府県の総人口 (Total_Population) のヒストグラム (10 階級) を示す。 2023 年度では東京 (1,409 万人) が最大、 鳥取 (54 万人) が最小で、 値域は 26.2 倍あり、 東京都 1 都だけで全人口の 11.3% を占める。 こうした歪んだ分布は分散処理の partition 設計に影響する。 住民 1 人 1 行のように県ごとのレコード数が人口に比例するデータを県で分割すると、 東京を含む partition が極端に重くなる (skew partition) ためだ。
分散処理で groupby('Prefecture') のような集約を行うとき、 partition (= データブロック) のサイズが大きく偏っていると、 一部の worker だけが長時間動き、 他は待機する「stragglers (落ちこぼれ)」問題が生じる。 SSDSE-B-2026 のような対数正規分布のデータでは、 partition 設計に**ハッシュベース** (Prefecture コードを hash して均等に振る) を用いることで、 各 worker が処理する行数を揃え、 全体のレイテンシを下げられる。 Dask では repartition(npartitions=N) や Spark では repartition(N, 'Prefecture') でこれを実現する。
分散処理の評価は「平均処理時間」だけでは不十分である。 SLA (Service Level Agreement, サービス品質契約) は通常 p95 (= 95 パーセンタイル) や p99 で定義されるため、 ばらつきこそが本質的に重要だ。 SSDSE-B-2026 (図 1 と同じ 47 行) を 10 万倍に複製した 470 万行のデータで、 Dask の並列度 (スレッドの worker 数) を 1, 2, 4, 8 と変えながら groupby 集計を各 100 回ずつ実施したときの実行時間分布を箱ひげ図で示す (環境で変わる目安)。
図 3 から、 並列度を 1 → 4 に上げても中央値は 1.44 倍 (135.5 → 93.8 ms)、 p99 は 1.45 倍 (144.1 → 99.1 ms) しか速くならず、 4 → 8 ではほぼ頭打ちになる (= worker を 4 倍にしても 4 倍速くはならない)。 1 台の中のスレッドで分けても、 処理の一部 (Python の GIL を持ったまま動く部分) は同時には進まず、 最後に partition ごとの集計結果をまとめる処理も残るためと考えられる。 この測定は 1 台の中で完結しているのでばらつきは小さいが、 複数台に分けると worker が増えるほど「いずれかの worker が GC・ネットワーク遅延・OS のスケジューリングに当たる確率」が上がり、 中央値より p99 が悪くなりやすい。 SLA 設計時は「中央値ではなく p99 で計画する」「タイムアウトとリトライを必ず実装する」が鉄則になる。
分散処理の性能評価には 2 つの軸がある。 Strong scaling は「データ量を固定して worker 数を増やす」評価で、 既存ジョブの高速化に対応する。 Weak scaling は「データ量と worker 数を比例して増やす」評価で、 データ増加に対するシステム拡張に対応する。 下表は、 SSDSE-B-2026 (47 行) を 100 万倍に複製した 4700 万行データで両方を測ったとしたらどうなるかを示す説明用の仮の数値で、 実測値ではない(このページでは 4700 万行の測定は行っていない。 効率の列は $E_s = T_1/(N T_N)$、 $E_w = T_1/T_N$ を仮の秒数から計算したもの)。
| worker 数 $N$ | Strong: データ 4700 万行固定 (秒) | Strong 効率 $E_s$ | Weak: データを $N$ 倍 (秒) | Weak 効率 $E_w$ |
|---|---|---|---|---|
| 1 | 38.2 | 1.000 | 38.2 | 1.000 |
| 2 | 20.4 | 0.936 | 40.1 | 0.953 |
| 4 | 11.6 | 0.823 | 42.7 | 0.895 |
| 8 | 7.3 | 0.654 | 46.9 | 0.814 |
| 16 | 5.1 | 0.468 | 52.4 | 0.729 |
この仮の例はAmdahl の法則 ($1 / (s + (1-s)/N)$, $s$ = 直列部分の割合) の典型例だ。 Strong scaling では worker 8 → 16 で効率が 65% → 47% に急落する、 というのは直列部分 (shuffle, scheduler 通信) が支配的になったときに典型的に現れる形である。 一方 Weak scaling では効率の落ち方が緩やかで、 「データ増加に応じてマシンを追加する」運用がより自然であることがわかる。 実務ではWeak scaling 優先で設計するのがほぼ常に正解になる。
分散処理のボトルネックはほぼ常にネットワーク帯域である。 SSDSE-B-2026 の 1 年度分 (47 行) を 1000 万倍に複製した 4.7 億行データを考え、 異なる集約操作の shuffle データ量を次の仮定から見積もった説明用の計算値を示す (実測値ではない)。 仮定: 使う数値列を 10 列 (float64 = 8 バイト) に絞り 1 行 80 バイト → 入力 = 4.7 億 × 80 B = 37.6 GB、 worker 数 25、 行は worker 間に均等に散らばる、 圧縮・キー列・シリアライズのオーバーヘッドは無視。 ここで shuffle 量とは「worker 間で送受信されたバイト総量」を指す。
| 操作 | shuffle 量 (GB) | 入力比 (%) | 理由 |
|---|---|---|---|
df.sum() (全体集計) | 約 0.000002 | ほぼ 0% | 各 worker が局所集計し、 10 列分の合計 (80 B) だけ送る: 25 × 80 B = 2 KB |
groupby('Prefecture').sum() | 約 0.0001 | 約 0.0003% | 各 worker が 47 グループの部分和を送る: 25 × 47 × 10 列 × 8 B = 94 KB |
sort_values('Total_Population') | 約 36.1 | 96% | 全データを値の範囲で再配置。 自分の担当範囲に元からある 1/25 以外が移動: 37.6 × 24/25 |
merge(big_df, big_df) | 約 72.2 | 192% | 両側全データを join key で再分散: 2 × 36.1 |
rolling(window=7).mean() | 約 0.00001 | ほぼ 0% | 隣接 partition の境界 6 行だけ転送: 24 境界 × 6 行 × 80 B ≈ 11.5 KB |
この表から、 「分散環境では sort と join を可能な限り避ける」「事前に partition key を揃えておく (= broadcast join, co-partitioned join)」が常識である理由がわかる。 1 Gbps ネットワーク (= 約 125 MB/s) で 72.2 GB を転送すれば、 72.2 / 0.125 ≈ 578 秒で理論最小でも約 10 分かかる。 これが Spark の broadcast() ヒントや Dask の set_index() 事前実行が推奨される所以だ。
分散処理は時間を買う行為であり、 必ずクラウド利用料という形で価格がつく。 SSDSE-B-2026 を 1 億行に拡張した処理 (60 GB) を各クラウドの代表的なマネージドサービスで実行した場合のコスト構造を比べるための説明用の仮の数値を示す (実測値でも各社の現行価格でもない)。 クラウドの料金は時期・リージョン・契約プラン・無料枠で変わるため、 実際の金額は必ず各社の最新の料金ページと料金計算ツールで確認すること。 ここで見てほしいのは「スキャン量課金の SQL エンジンは短時間・少額、 クラスタ課金の Spark 系は起動時間分の費用がかかる」という構造の違いである。
| サービス | 構成 | 実行時間 | 仮の概算コスト (USD、 説明用) | 向くワークロード |
|---|---|---|---|---|
| AWS EMR (Spark) | m5.xlarge × 4 + master | 12 分 | $0.42 | 定常 ETL バッチ |
| AWS Glue | G.1X × 4 DPU | 15 分 | $0.66 | サーバレス ETL |
| GCP Dataproc Serverless | 4 vCPU × 4 executor | 11 分 | $0.38 | アドホック Spark |
| GCP BigQuery | on-demand (60 GB scan) | 8 秒 | $0.30 | SQL 集計 |
| Azure Synapse (Spark pool) | Small × 4 node | 14 分 | $0.52 | Microsoft 環境統合 |
| Databricks (AWS) | i3.xlarge × 4 | 9 分 | $0.71 | 高頻度開発、 Delta Lake |
SQL で解ける問題なら BigQuery (= MPP query engine) が最速・最安、 複雑な ETL や ML 前処理なら Spark 系 (EMR, Dataproc, Databricks)、 サーバ管理を避けたければ Glue や Dataproc Serverless が向く、 という典型的なトレードオフがある。 SQL 集計だけが目的なら、 BigQuery のようなスキャン量課金のエンジンに load して SQL で集計するのが最も簡単で、 費用は「スキャンしたバイト数 × その時点の単価」で決まる (必要な列だけ SELECT し、 パーティション分割でスキャン量を減らすほど安くなる)。
大規模分散処理では「worker が途中で死ぬ」ことが日常的に起こる。 Spark や Dask は次の 2 方式で対処する。 (1) lineage 再計算: 各 partition の生成手順 (有向非巡回グラフ) を保持し、 失敗時は依存元から再計算する。 (2) checkpoint: 中間結果を分散ストレージ (HDFS, S3) に永続化する。 lineage は計算が短ければ安価、 checkpoint は計算が長く失敗ペナルティが大きい場合に有効である。
| 方式 | 復旧時間 | 追加コスト | 向く処理 |
|---|---|---|---|
| lineage のみ | 依存全部再計算 (高) | ほぼゼロ | 短い変換チェイン |
| checkpoint (S3) | 最後の checkpoint から (低) | ストレージ I/O + 容量 | 反復計算 (ML, graph) |
| replication (= broadcast) | 即座 (= レプリカから) | $N$ 倍のメモリ | 小サイズ参照テーブル |
| speculative execution | 予測的に複製実行 | CPU 余剰 + 重複計算 | straggler 多発環境 |
SSDSE-B-2026 のような小データの分散処理 (= ベンチマークや教育用途) では失敗確率が低いため lineage 任せで十分だが、 数百 TB 規模のジョブでは「30 分ごとに checkpoint」「join 直前に必ず checkpoint」のような運用ルールが必須となる。 Spark の場合 df.checkpoint() もしくは df.persist(StorageLevel.DISK_ONLY), Dask では client.persist(df) で明示する。
ここまでの内容を踏まえ、 次の 6 問に答えてみよう。 すべて SSDSE-B-2026 や本ページに出てきた数値・概念に基づく。 答えは各問の下に折りたたんでいる。
Q1. SSDSE-B-2026 (47 行) で groupby を実行するとき、 pandas と PySpark local モードでは中央値で何倍の差があるか? 本ページの表(説明用の仮の数値)で答えよ。
表の仮の値では pandas が 0.18 ms, PySpark local が 92.1 ms なので、 約 511 倍 pandas が速い。 倍率そのものは環境で変わるが、 SSDSE-B-2026 規模では分散処理が逆効果になるという向きは変わらない。
Q2. 図 1 で示した、 pandas と Dask local の処理時間が交差する行数の目安は何行か?
図 1 の測定では交差点は現れなかった。 Dask / pandas の時間の比は 47 行で約 66 倍、 47 万行で 1.84 倍、 1,410 万行で 1.09 倍と縮まるが、 メモリに収まる範囲では最後まで pandas が速い。 Dask が効くのは 1 台のメモリに載らない量を扱うときである。
Q3. Strong scaling と Weak scaling の違いを 1 行で説明し、 実務でどちらを優先するべきか答えよ。
Strong = データ固定で worker を増やす、 Weak = データと worker を比例して増やす。 実務ではWeak 優先。 Strong は Amdahl の法則で頭打ちになり、 worker を増やしても効率が下がる。
Q4. 4.7 億行データ (1 行 = float64 10 列 = 80 バイト) を 25 worker に均等に分けて sort_values を実行した場合、 shuffle 量は約何 GB になり、 入力比でいうと何 % か? (行の行き先は 25 worker に均等に散らばると仮定する)
約 36.1 GB (入力比 96%)。 手順: (1) 入力 = 4.7 × 108 行 × 80 B = 37.6 GB。 (2) sort では各行が値の範囲で決まる担当 worker へ送られ、 行き先が今いる worker と同じになる確率は 1/25 なので、 移動する割合は 24/25 = 0.96。 (3) shuffle 量 = 37.6 × 0.96 ≈ 36.1 GB。 worker 数を増やすほど移動割合 (N−1)/N は 100% に近づき、 shuffle 量はほぼ入力サイズと等しくなる。 分散環境で sort は最も高価な操作のひとつ。 (実際の量は圧縮やシリアライズ形式で変わる。)
Q5. SSDSE-B-2026 を 60 GB に拡張した場合、 SQL 集計だけが目的なら最も安価で速いマネージドサービスは何か?
GCP BigQuery のようなサーバレスの SQL エンジン (列指向ストレージ + MPP) が候補になる。 クラスタを起動・維持する時間がかからず、 費用はスキャン量に比例する (on-demand の場合)。 金額は料金改定・リージョン・プランで変わるので、 上のコスト表の数値は説明用の仮の値として扱い、 実際の単価は公式の料金ページで確認する。 必要な列だけを読むとスキャン量 (= 費用) が減る。
Q6. 反復計算 (= ML の勾配降下を 100 epoch 回すなど) を分散環境で行うとき、 lineage 任せでは何が問題になるか? 解決策と合わせて答えよ。
lineage chain が深くなり、 worker 失敗時の再計算コストが膨大になる。 解決策は checkpoint を定期的に取る (Spark df.checkpoint() / Dask client.persist())。 反復ごと、 もしくは数イテレーションごとに中間結果を S3/HDFS に書き出す。
分散処理を理解する上で避けて通れないのが CAP 定理 (Brewer の定理, 2000 年提唱, 2002 年 Gilbert-Lynch 形式証明) である。 これは「Consistency (一貫性)」「Availability (可用性)」「Partition tolerance (分断耐性)」の 3 つを同時には満たせない、 という分散システム設計の根源的制約を述べた定理だ。 ネットワーク分断は現実世界では必ず起こるため (= P は捨てられない)、 実務的な選択は「C か A か」の二択になる。 SSDSE-B-2026 を例にとって考えてみよう。 もし全 47 都道府県のデータを 3 つの worker に分散保持し、 ある worker がネットワーク障害で他から切り離されたとき (= partition)、 設計者は次の選択を迫られる。 (1) C 優先 (CP 系): 切り離された worker は読み書きを拒否 (= unavailable) し、 全体の一貫性を守る。 HBase, MongoDB (デフォルト), ZooKeeper などが該当。 (2) A 優先 (AP 系): 切り離された worker も独立に読み書きを継続し、 後で整合性を解決 (= eventual consistency) する。 Cassandra, DynamoDB, CouchDB が該当。
分散処理の用途別に整理すると、 ETL バッチ処理 (= SSDSE-B-2026 の年次集計のような長時間ジョブ) は CP 系で良い (= ジョブが多少遅れても結果が正しい方が重要)。 一方、 オンライン推薦エンジンや IoT センサ収集 (= 1 秒の停止が損失を生む) は AP 系を選ぶ。 さらに新しい潮流として PACELC 定理 (2010 年, Abadi 提唱) がある。 これは「分断がないとき (= Else) も、 Latency (遅延) と Consistency のトレードオフがある」という拡張で、 例えば同期レプリケーション (= 強整合, 高遅延) vs 非同期レプリケーション (= 弱整合, 低遅延) の選択が常に発生することを明文化している。 PACELC で分類すると、 ほとんどの実用システムは「PA/EL」(分断時は A, 通常時は低遅延優先 = 弱整合) になっている。 SSDSE-B-2026 を BigQuery に積んで集計するケースは「PC/EC」(常に強整合) であり、 これが BigQuery が「分析用途」に振り切られた設計である理由でもある。
分散処理エンジン (Spark, Dask, Flink) は単独では動かない。 必ず「データをどう保存するか」を決めるストレージ層とセットになる。 ここでは SSDSE-B-2026 のような表形式データを分散ストレージに置く際の典型的な選択肢を整理する。 まず物理ストレージとして 2 系統がある。 (1) HDFS (Hadoop Distributed File System): 各 worker のローカルディスクを束ねて 1 つの巨大ファイルシステムに見せる。 データ局所性 (= 計算を data の近くに送る) を優先するが、 ストレージとコンピュートが結合しているため拡張が難しい。 (2) オブジェクトストレージ (S3, GCS, Azure Blob): ストレージとコンピュートを分離。 worker は計算するときだけ起動し、 終わったら破棄できるため、 spot/preemptible 活用と相性がよく、 現代の標準である。 SSDSE-B-2026 のような数 MB のデータを S3 に置く場合、 ストレージコストはほぼゼロ (約 0.023 ドル/GB/月) で、 むしろ API 呼び出し回数 (約 0.0004 ドル/1000 requests) の方が気になる規模になる。
ファイル形式の選択も性能に大きく効く。 SSDSE-B-2026 を CSV のまま置くと、 1 行ずつ全列読み込む必要があり (= 行指向)、 「総人口の県別合計だけ欲しい」というクエリでも全列を読まされる。 これを Parquet (Apache Parquet, 列指向 + 圧縮) に変換すると、 必要な列だけ読めるようになり (= column pruning)、 同じクエリの I/O が 1/20 以下になることもある。 さらに Delta Lake (Databricks 開発) や Apache Iceberg (Netflix 開発) を使うと、 ACID トランザクション (= 同時書き込み時の一貫性保証), time travel (= 過去版へのロールバック), schema evolution (= 列追加・型変更の追跡) といったデータベース的機能が分散ファイルでも使えるようになる。 SSDSE-B-2026 のように年次更新される公的データは、 「2024 年版」「2025 年版」「2026 年版」を Delta Lake のバージョン管理で持っておくと、 「2024 年版で集計した結果を再現したい」というニーズに即座に応えられる。
これまで本ページで扱ってきたのは「バッチ処理」、 つまり「データを溜めてから一括処理する」モードだった。 SSDSE-B-2026 のような年次統計はこのモデルにぴったり合う。 しかし、 IoT センサ, クリックストリーム, 金融市場データのように「データが連続的に流れ込む」場合はストリーム処理が必要になる。 分散ストリーム処理エンジン (Apache Flink, Spark Structured Streaming, Apache Beam, Kafka Streams) は、 (1) ミリ秒〜秒オーダーの低遅延で連続データに対する集約・結合・パターン検出を行う, (2) windowing (= 5 分窓, 1 時間窓などの時間範囲集計) を扱う, (3) late-arriving data (= 遅れて届くデータ) を watermark で処理する、 といった機能を備える。 SSDSE-B-2026 を例にすると、 「全国の出生件数を 1 時間ごとに集計する」「県別の事故発生件数を 5 分窓で監視する」といったユースケースが該当する (現実には公的統計はストリーム化されていないため、 これは思考実験)。
過去にはバッチとストリームを別系統で持つ Lambda アーキテクチャ (Nathan Marz 提唱) が主流だった。 これは「正確だが遅いバッチ層 + 速いが近似のストリーム層を並行運用し、 サービング層で結果を合成する」というもの。 SSDSE-B-2026 のような正確な集計には Hadoop/Spark バッチ、 リアルタイム速報には Storm/Flink、 という二重実装が必要だった。 これに対し、 Jay Kreps (Confluent 創業者) が 2014 年に提唱した Kappa アーキテクチャは「全てをストリーム処理として扱い、 バッチは『長い時間範囲のストリーム』として再処理する」という単一化案で、 Kafka + Flink で実現される。 さらに近年 (= 2020 年代) は Unified Batch and Streaming (Apache Beam, Spark Structured Streaming) として、 同じ API でバッチ・ストリーム両方を書けるエンジンが主流化しており、 「バッチか, ストリームか」を実装後に切り替えられるようになっている。 SSDSE-B-2026 規模では普通バッチで十分だが、 大規模システムを設計するときはこの 3 アーキテクチャの違いを理解しておく必要がある。
分散処理の古典的な金言に「Move computation to data, not data to computation (計算をデータに送れ、 データを計算に送るな)」がある。 これは Hadoop の設計哲学で、 「数 TB のデータを集計するために数 KB のコードを worker に送る方が、 数 TB のデータを集約マシンに転送するより圧倒的に効率的」という観察に基づく。 SSDSE-B-2026 (数 MB) のような小データではこの議論は無意味だが、 ペタバイト級データではデータ転送に物理的限界 (= 1 Gbps ネットワークで 1 PB を転送するには 90 日以上) があるため、 計算側を動かす方が常に正解になる。 Hadoop, Spark on YARN/Kubernetes, Dask にはlocality-aware scheduling が組み込まれており、 タスクは「データのある worker」を優先して割り当てられる。 もし優先 worker が満杯なら、 同じラック内の別 worker (= rack-local), 最後に異なるラック (= any) と落としていく。
しかしクラウド時代になると、 ストレージとコンピュートが物理的に分離した (= S3 のデータを EC2 から読む) ため、 真の data locality は失われた。 代わりに「キャッシュ局所性」と「partition affinity」が重要になる。 例えば SSDSE-B-2026 を 100 倍に複製したデータを S3 に置き、 Spark で cache() や persist() を使うと、 一度読んだ partition は worker のローカルメモリ/ディスクに保持され、 2 回目以降のクエリは S3 を経由しない (= ローカル読み)。 これにより S3 → EC2 のネットワーク帯域 (約 10-25 Gbps per node) を消費せず、 ローカル NVMe (約 6 GB/s) の速度で読める。 同様に、 partition key を揃えた join (= co-partitioned join) では、 worker 間転送が不要になる。 SSDSE-B-2026 の集計を Prefecture でパーティション化しておけば、 「人口 ⨝ 出生数」のような結合がほぼゼロコストになる。 これらは分散処理の最重要の高速化技法であり、 SSDSE-B-2026 規模のテストデータで挙動を理解しておくと、 本番の大規模データで効きが圧倒的に違ってくる。
本ページの締めくくりとして、 実務で分散処理ジョブを本番投入する前に必ず確認すべき項目を 10 個リストアップする。 SSDSE-B-2026 のような小データではほぼ無視できる項目もあるが、 本番データに拡張するときは全て検討しておくべき:
reduceByKey を groupByKey + reduce より優先。 broadcast 結合を可能なら活用。--maximum_bytes_billed, EMR の job timeout 等を必ず設定。特に最後の「テストデータでの dry-run」が重要だ。 ペタバイトデータでの初実行は数時間〜数日かかり、 失敗するとそのまま課金される。 SSDSE-B-2026 のような小さな実データでロジックを完全に検証してから本番データに展開する、 という流れが事故を防ぐ唯一の道筋になる。 「分散処理の練習は小データで」「本番は十分にテストしてから」を肝に銘じよう。 SSDSE-B-2026 は無料で取得でき, 47 行という扱いやすいサイズで全てのテストパターンを高速回転できる, 学習と本番リハーサルの両方に使える稀有な公的データセットである。
分散処理基盤を組織で共有運用する場合、 セキュリティとガバナンスが新たな課題として浮上する。 SSDSE-B-2026 のような公開データなら気にしなくてよいが、 個人情報・取引情報・医療情報を扱う場合は次の論点が常に伴う。 (1) アクセス制御: どの worker がどのデータを読めるか。 Hadoop は Kerberos + Ranger, Spark on Databricks は Unity Catalog, AWS は IAM + Lake Formation で行・列・タグ単位のアクセス制御を実現する。 (2) 暗号化: at-rest (= 保存時) と in-transit (= 転送時) の両方が必要。 S3 SSE-KMS, HDFS Transparent Encryption, TLS 1.3 などが標準。 (3) 監査ログ: 誰がいつ何を読んだか。 CloudTrail, AWS GuardDuty, Apache Atlas でログを集中管理し、 異常アクセスを検知する。 (4) data lineage: ある出力データがどの入力から生成されたか。 OpenLineage, Marquez, dbt が代表的なツールで、 GDPR の「削除権」対応にも使える。 SSDSE-B-2026 のような公開統計を扱う段階ではこれらは過剰だが、 自治体の住民データや企業の顧客データに移行するときは最初から設計に組み込むのが鉄則。
特にマルチテナント環境 (= 複数チーム・複数プロジェクトが同じクラスタを共有) では、 「ジョブごとの分離 (= isolation)」と「リソース割当 (= quota)」が重要になる。 YARN や Kubernetes は namespace + resource quota でこれを実現する。 例えば「データ分析チームは CPU 100 コア・メモリ 500GB まで、 ML 推論チームは GPU 4 枚まで」のように上限を切ることで、 一つの暴走ジョブが他テナントを巻き込まないように防御する。 SSDSE-B-2026 のような数 MB データを扱う限り resource quota はほぼ無意味だが、 本番運用では「コスト管理 × セキュリティ × SLA」の三本柱として必須になる。 分散処理を学ぶときは、 まず小データでロジックを習得し、 次に小クラスタ (= local cluster) で性能特性を理解し、 最後にマルチテナント本番環境で運用 (= governance) を学ぶ、 という段階的なアプローチが推奨される。
分散処理エンジン (Spark, Presto/Trino, BigQuery, Redshift, Snowflake, ClickHouse など) の性能比較には、 業界標準のベンチマークスイートが用いられる。 代表は TPC-DS (Transaction Processing Performance Council Decision Support, 99 クエリ, 小売業をモデル化), TPC-H (22 クエリ, 卸売をモデル化), SSB (Star Schema Benchmark, 13 クエリ, スター型データウェアハウス)。 これらは数 GB 〜 数 TB のサイズで実行され、 クエリ完了時間, 同時実行性能, データロード時間, ストレージ効率などを総合的に評価する。 SSDSE-B-2026 は本格的なベンチマーク対象ではないが、 「47 都道府県 × 約 100 列」というスター型に近い構造を持つため、 「Prefecture を dimension, 各統計指標を fact とみなす」という発想で TPC-DS 的なクエリ (例: 「地方ブロック別に年代別人口を集計し、 過去 10 年の伸びでソートする」) を書く練習ができる。 ベンチマークで上位を取るシステムは、 (1) 列指向ストレージ + 圧縮の効率, (2) クエリ最適化器 (cost-based optimizer) の賢さ, (3) ベクトル化実行エンジン, (4) 並列度の動的調整、 という共通の強みを持つ。 SSDSE-B-2026 を題材に「Spark で書いた集計クエリ」を「DuckDB で書き換える」「BigQuery に積み替える」と比較すると、 数百倍の速度差が生まれる場面を体験でき、 「適切なエンジン選び」の重要性を肌で理解できる。
分散処理ジョブが「思ったより遅い」「メモリ不足で落ちる」「コストが想定の 3 倍」となる典型的原因と、 対策は次の通りである。 SSDSE-B-2026 規模では問題にならなくても、 本番では必ず効いてくる定石なので、 小データで挙動を確認しておくと良い:
year=2026/prefecture=Tokyo/ のようにディレクトリ階層で分割して保存すると、 WHERE 句で関係ない partition を完全にスキップできる。 SSDSE-B-2026 の年次データを year=YYYY で分割しておくのは将来拡張のための基本。broadcast(df) ヒントで強制可能。 SSDSE-B-2026 を fact テーブルとし、 県コード変換マスタを broadcast するのが典型例。cache()。 ただし不要になったら unpersist() でメモリ解放。 SSDSE-B-2026 で 5 回参照する中間集計があれば cache、 1 回だけならしない。repartition (shuffle 発生), 減らすだけなら coalesce (shuffle なし) を使う。 ファイル書き出し前は coalesce(1) で 1 ファイル化することが多いが、 大データではこれが OOM の原因になる。col / 1000 の方が桁違いに速い。spark.sql.adaptive.enabled=true で有効化。 SSDSE-B-2026 のような小データでは効果薄だが、 本番では必ず有効にする。これら 7 項目を意識するだけで、 同じハードウェア・同じデータでも処理時間が 5-20 倍改善することは珍しくない。 「分散処理は遅い」と感じている人の多くは、 実は分散処理の使い方を最適化していないだけであり、 上記の定石を順番に適用していくと劇的に変わる。 SSDSE-B-2026 のような小データで各定石の効果を確認しながら学ぶのが、 遠回りのようで最短ルートになる。
現代の分散処理を理解するには、 その歴史的経緯を押さえることが理解の近道になる。 2003 年に Google が GFS (Google File System) 論文を, 翌 2004 年に MapReduce 論文を発表したのが起点。 Yahoo! の Doug Cutting が両者を OSS 実装した Apache Hadoop (HDFS + MapReduce) が 2006 年に登場し、 「コモディティハードウェアでペタバイトを扱う」という常識が世界に広まった。 SSDSE-B-2026 のような小データには大袈裟だが、 「データを分割し, 各 worker で map 処理し, key 単位で reduce する」というモデル自体は今も生きており、 Spark の groupBy().agg() の内部実装はまさにこれである。
2010 年代に入ると, MapReduce の「ディスク I/O が多すぎて遅い」「機械学習の反復計算に不向き」という弱点を克服する Apache Spark (UC Berkeley, 2010 年) が登場。 in-memory 計算と RDD (Resilient Distributed Dataset) の lineage 機構で 10-100 倍の高速化を実現した。 SSDSE-B-2026 の集計を Hadoop MapReduce で書くと数百行になるが、 Spark なら 3-5 行で済む。 さらに 2014 年に Spark DataFrame API が導入され, pandas ライクな書き味で分散処理ができるようになった。 並行して, Google は Dremel (2010) を商用化した BigQuery (2011) を一般公開し, 「SQL だけでペタバイトを 10 秒で集計する」世界を見せた。 Snowflake (2014 年公開) や Databricks Lakehouse (2020 年代) はこの流れの上にあり, ストレージとコンピュート分離の現代的な姿を完成させた。 SSDSE-B-2026 を BigQuery にロードして集計する経験は, この 20 年の進歩の最終形を数百円で体験することに等しい。
2020 年代に入ると, Lakehouse (Delta Lake, Apache Iceberg, Apache Hudi) の概念が普及し, 「データウェアハウス (構造化, 高速 SQL) とデータレイク (非構造化, 安価) の良いとこ取り」が標準アーキテクチャとなった。 さらに最近では クエリ高速化エンジン として DuckDB (2019 年, in-process OLAP), Polars (2021 年, Rust 製 DataFrame) が台頭し, 「単一マシンで数十 GB を秒で集計する」ことが個人 PC で可能になった。 これは「分散処理が必要なサイズの閾値が年々上がっている」ことを意味する。 2010 年代は「10 GB 超えたら分散」だったのが, 2026 年現在は「100 GB 超えたら分散」へと事実上シフトしている。 SSDSE-B-2026 を 100 倍, 1000 倍に複製したデータを Polars/DuckDB で集計すると, Spark を立ち上げるよりも速く終わる, という現代的な感覚を体験することを強く勧める。 「Spark を学ばないと分散処理が分からない」というのは半分正しく半分間違いで, 現代では「まず Polars/DuckDB で限界まで頑張り, それでも足りなければ Spark/Dask」というのが現実解になっている。
分散処理は「データが単一マシン RAM に収まらなくなったとき」「処理時間が SLA を超えるとき」「リアルタイム取り込み量が単一マシン I/O を超えるとき」のいずれかを満たして初めて検討する。 SSDSE-B-2026 (47 行 / 数 KB) の段階では pandas が常勝であり、 分散処理は逆効果になる。 本ページの図 1-3 と表 (実測値、 shuffle 量、 クラウドコスト、 fault tolerance) は、 この判断を**実データの数値で支える**ためのリファレンスである。 「とりあえず Spark」は禁句。 「データサイズ → 候補システム → 実測 → 採否」の順で必ず判断しよう。 この判断手順を身につけることが, 分散処理を扱う上で最も重要なスキルになる。 ツールの細部より, 「いつ使う / いつ使わない」を即答できるエンジニアになることを目指してほしい。 これが本ページのもっとも大事な持ち帰りポイントである。
本ページで学んだ内容を一度きりで終わらせず, 実際の手で動かして確かめることを強く推奨する。 SSDSE-B-2026 は無料で入手できる公的統計データであり, pandas での集計, Dask local cluster での再現, BigQuery への load と SQL クエリ, さらには複製サイズを段階的に増やしたスケーリング実験まで, 全てを自宅 PC とクラウド無料枠で実行できる。 こうした実験を通じて, 「自分のマシンで処理時間がどう変化するか」「partition 数を変えると何が起こるか」「列指向ファイル形式 (Parquet) はどれだけ速くなるか」を体感すると, 教科書的な知識が一気に実務知に転化する。 また, クラウドコストの実体験 (例: BigQuery のクエリ実行前に表示される「処理されるバイト数」の見積もりを見て、 スキャン量に比例して料金が決まることを確かめる) を一度経験すると, 「思いつきで巨大クエリを叩く」ことの怖さも理解できる。 分散処理は「動かす」ことが最大の学びになる分野なので, 本ページを参照しながら手を動かしてほしい。 SSDSE-B-2026 だけで, MapReduce 以来 20 年の分散処理の進化を一通り再現体験できる, 稀有な教材になっている。
🍰 まずはやさしく
チームで作業を分担するようなものです。
待ち時間を減らして効率よく計算するために使います。
1人で全国を回るより、47人で分担する方が早いです。
この章では計算を分ける3つのやり方について読みます。
「47 都道府県を 1 人で巡業」 vs 「47 人が同時に各県で集計し東京に集約」のどちらが早いか考えてください。後者がデータ分散処理の発想です。集計内容 (map)、集約方法 (reduce)、移動コスト (シャッフル) を慎重に設計しないと、人数を増やしても遅くなることもあります。
3 種類の並列性:
| 種類 | 分けるもの | 例 |
|---|---|---|
| データ並列 | 入力データを分割 | 47 都道府県別に同じ集計を分担 |
| タスク並列 | 違う計算を同時に | 人口集計と高齢化率集計を並行 |
| パイプライン並列 | 段階を流す | 読込→前処理→推論→保存を流れ作業 |
MapReduce のメタファー: map は「各都道府県で集計する係」、reduce は「47 県の結果を 1 つにまとめる係」、shuffle は「集計表を持ち寄って同じ項目どうしを並べる作業」。
🍰 まずはやさしく
計算の速さを測るためのルールです。
どれくらい効率的に早くなったかを知るために使います。
部活の練習メニューを分担して時間を短くする例に似ています。
この章では高速化の限界や計算式について読みます。
Speedup (高速化倍率):
$$ S(n) = \frac{T(1)}{T(n)} $$
$T(n)$ は $n$ ノードで処理したときの実時間。理想は $S(n) = n$。
Amdahl の法則:並列化できる部分の比 $p$
$$ S(n) = \frac{1}{(1-p) + p/n} \xrightarrow{n \to \infty} \frac{1}{1-p} $$
どんなに $n$ を増やしても、逐次部分 $1-p$ が高速化の限界を決める。
Gustafson の法則:問題サイズを増やせるなら
$$ S(n) = (1-p) + p \cdot n $$
MapReduce の意味論:
$$ \mathrm{map}: (k_1, v_1) \to [(k_2, v_2)], \quad \mathrm{reduce}: (k_2, [v_2]) \to (k_2, v_3) $$
SSDSE では $k_1$=行番号、$v_1$=行、$k_2$=都道府県、$v_2$=値リスト、$v_3$=合計や平均などの集計結果。
CAP 定理: 分散システムは Consistency / Availability / Partition tolerance のうち同時に満たせるのは 2 つまで。
$$ C \cap A \cap P = \emptyset $$
| 記号 | 読み方 | 意味・例 |
|---|---|---|
| $n$ | number of workers | ノード数 (47 都道府県なら 47) |
| $T(n)$ | time | $n$ 並列での実行時間 |
| $p$ | parallelizable fraction | 並列化可能比率 (0〜1) |
| $E(n)$ | efficiency | $S(n)/n$。100%が理想 |
| $k_1, k_2$ | keys | map/reduce のキー (例: 都道府県) |
| $v_1, v_2, v_3$ | values | 入力 / 中間 / 集計後の値 |
SSDSE-B-2026 は 47 都道府県 × 12 年 (2012-2023) = 564 行。これは「47 個に綺麗に切れるデータ」のおもちゃ例として最適です。
Amdahl の法則の数値例:
| 並列化率 $p$ | $n=4$ の Speedup | $n=47$ の Speedup | 上限 |
|---|---|---|---|
| 0.50 | 1.60× | 1.96× | 2.00× |
| 0.90 | 3.08× | 8.32× | 10.00× |
| 0.99 | 3.88× | 32.7× | 100× |
「ほぼ並列化できる ($p=0.99$) 処理」でも、47 並列で 32 倍にしかならない。逐次部分の最適化が命。
SSDSE 集計のコスト見積もり:
| 処理 | 1 都道府県あたり | 47 都道府県逐次 | 47 並列 |
|---|---|---|---|
| 行平均集計 | 1 ms | 47 ms | 3 ms (起動オーバーヘッド込み) |
| ML 特徴量計算 | 300 ms | 14.1 s | 450 ms |
| 回帰モデル fit | 800 ms | 37.6 s | 1.1 s |
合成データで並列化率 70% のジョブを N ノードで実行したときの速度向上を計算する。
| N | S(N) |
|---|---|
| 1 | 1.00 |
| 2 | 1.54 |
| 4 | 2.11 |
| 8 | 2.58 |
| ∞ | 3.33 |
1 2 3 4 5 6 | import numpy as np P = 0.70 N = np.array([1, 2, 4, 8]) S = 1 / ((1-P) + P/N) print(f"S(N): {S.round(3)}") print(f"上限 (N→∞): {1/(1-P):.3f}") |
💬 手計算 (Step 3) S(4)≈2.105 と Python 出力が完全一致。
① まずは逐次処理 (ベースライン)
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 | # Total_population / Income_per_capita / Aging_rate という列は SSDSE-B-2026 に無い。 # A1101(総人口)・L3221(消費支出)に置き換え、高齢化率は派生させる。 import pandas as pd import time df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=[1]) df['高齢化率'] = df['A1303'] / df['A1101'] * 100 def heavy_aggregate(group): # 何か重い計算 (例: 移動平均・分位点・線形回帰) return { 'mean_pop': group['A1101'].mean(), 'mean_inc': group['L3221'].mean(), 'aging': group['高齢化率'].mean(), 'rows': len(group), } t0 = time.time() result_seq = {p: heavy_aggregate(g) for p, g in df.groupby('Prefecture')} print(f'逐次 : {time.time() - t0:.3f} s, 件数 = {len(result_seq)}') |
💬 47 都道府県ぶんの集計が 0.003 秒で終わった。1 グループはわずか 12 行(12 年度分)なので、ここで重いのは計算ではなく groupby の準備と辞書の作成で、並列化の出番はない。分散処理が効いてくるのは、1 グループの計算に秒単位かかるか、データがメモリに載らない規模になってからである。
② concurrent.futures.ProcessPoolExecutor で 47 並列
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 | # ※ このブロックはノートブックやブラウザでは動かない。 # ProcessPoolExecutor は別プロセスに関数を送るため、その関数が # インポート可能な場所(.py ファイル)に無いと pickle できない。 # 手元で試すときは、この内容を worker.py などに保存して python worker.py で動かす。 # ここでは同じ処理をスレッドで行い、結果が一致することだけ確かめる。 from concurrent.futures import ThreadPoolExecutor import time def task(args): pref, group = args return pref, heavy_aggregate(group) groups = list(df.groupby('Prefecture')) t0 = time.time() with ThreadPoolExecutor(max_workers=8) as ex: result_par = dict(ex.map(task, groups)) print(f'並列 : {time.time() - t0:.3f} s') print('逐次と同じ結果か:', result_par == result_seq) |
💬 スレッド 8 本で 0.004 秒と、逐次の 0.003 秒より遅くなった。pandas の小さな集計は GIL を握ったまま進むのでスレッドでは並列に走らず、スレッドを作って結果を集める手間の分だけ時間が増える。結果の辞書が逐次と一致した(True)ことは、並列化で順序が入れ替わっても集計値が変わらないことの確認として意味がある。
③ Dask: 大規模 CSV を chunk 並列
1 2 3 4 5 6 7 8 9 10 11 12 | # Total_population / Income_per_capita / Aging_rate という列は SSDSE-B-2026 に無い。 # A1101(総人口)・L3221(消費支出)に置き換え、高齢化率は派生させる。 import dask.dataframe as dd ddf = dd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=[1], blocksize='8MB') ddf['高齢化率'] = ddf['A1303'] / ddf['A1101'] * 100 agg = ddf.groupby('Prefecture').agg({ 'A1101': 'mean', 'L3221': 'mean', '高齢化率': 'mean', }).compute() # ここで初めて実計算 (遅延評価) print(agg.head()) |
💬 北海道の A1101 平均 529.7 万人は、2012〜2023 年度の 12 年分(546.5 万人から 509.2 万人へ減少)を平均した値で、2023 年度の人口ではない。dask は compute() を呼ぶまで計算せず、約 36 万バイトの CSV は blocksize='8MB' では 1 パーティションにしかならないので、ここでの結果は pandas と同じ値を遠回りして得ているだけ。パーティションが 1 つかどうかは ddf.npartitions で確かめられる。
④ MapReduce 風に 47 県を手で書く
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 | # Total_population / Income_per_capita / Aging_rate という列は SSDSE-B-2026 に無い。 # A1101(総人口)・L3221(消費支出)に置き換え、高齢化率は派生させる。 from collections import defaultdict # map: (行) -> (都道府県, 値) の列を吐く def mapper(row): yield row['Prefecture'], row['A1101'] # shuffle: 同じキーを集める buckets = defaultdict(list) for _, row in df.iterrows(): for k, v in mapper(row): buckets[k].append(v) # reduce: 各キーで平均(合計 ÷ 件数 = 12 年度分の平均人口) result = {k: sum(vs) / len(vs) for k, vs in buckets.items()} print(list(result.items())[:5]) |
💬 map → shuffle → reduce を手で書いた結果、北海道 5,297,195.58、青森県 1,271,937.42 と、直前の dask の A1101 平均と小数点以下まで一致した。reduce で平均を出すときは (合計, 件数) の組を持ち回って最後に割るのが MapReduce の定石で、途中の平均どうしを平均するとグループの件数が違う場合に値がずれる。iterrows は 564 行でも遅い部類なので、この書き方は仕組みの説明用と割り切る。
⑤ Ray でモデル学習を 47 並列
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 | # Total_population / Income_per_capita / Aging_rate という列は SSDSE-B-2026 に無い。 # A1101(総人口)・L3221(消費支出)に置き換え、高齢化率は派生させる。 df['高齢化率'] = df['A1303'] / df['A1101'] * 100 import ray from sklearn.linear_model import LinearRegression ray.init(ignore_reinit_error=True) @ray.remote def fit_one(pref, sub): X = sub[['A1101', '高齢化率']].values y = sub['L3221'].values if len(sub) < 3: return pref, None m = LinearRegression().fit(X, y) return pref, (m.coef_.tolist(), m.intercept_) futures = [fit_one.remote(p, g) for p, g in df.groupby('Prefecture')] models = dict(ray.get(futures)) print(list(models.items())[:3]) |
⑥ Amdahl の法則を可視化
1 2 3 4 5 6 7 8 9 | import numpy as np import matplotlib.pyplot as plt n = np.arange(1, 48) for p in [0.5, 0.9, 0.99]: plt.plot(n, 1 / ((1 - p) + p / n), label=f'p={p}') plt.xlabel('# workers (都道府県数)'); plt.ylabel('Speedup') plt.legend(); plt.grid(True) plt.title('Amdahl: 47 都道府県を並列にしても逐次部が上限を決める') |
🎯 このコードでやること:SSDSE-B-2026 の 47 都道府県 × 12 年分 (2012-2023) のデータを、 手作りの MapReduce で「都道府県平均人口」を集計する。 map / shuffle / reduce の 3 ステップを明示する教育目的。
📥 入力データ (SSDSE-B-2026.csv の先頭、 cp932 エンコーディング):
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 | import pandas as pd from collections import defaultdict # SSDSE-B-2026 を読み込み (cp932) df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=[1]) df.columns = ['Year','Code','Pref','Pop'] + list(df.columns[4:]) # === MAP ステップ === def mapper(row): """(都道府県, 人口) のキーバリューペアを emit""" yield (row['Pref'], int(row['Pop'])) # === SHUFFLE ステップ === buckets = defaultdict(list) for _, row in df.iterrows(): for k, v in mapper(row): buckets[k].append(v) # === REDUCE ステップ === def reducer(key, values): return key, sum(values) / len(values) result = dict(reducer(k, vs) for k, vs in buckets.items()) # 結果を人口降順で表示 for pref, mean_pop in sorted(result.items(), key=lambda x: -x[1])[:5]: print(f'{pref}: 平均人口 {mean_pop:,.0f}') |
📤 実行例 (実際の出力):
💬 結果の読み方:47 都道府県 × 12 年 (2012-2023) のデータが、 都道府県ごとに集約された (47 行に縮約)。 これが MapReduce の本質 — 「巨大データをキー単位で集約することで、 47 並列のシャッフル + 集約が成立する」。 東京が約 1375 万、 鳥取が約 56 万で 24 倍の偏りも確認できる。
🎯 このコードでやること:47 都道府県を multiprocessing.Pool で並列処理し、 逐次実行と比較する。 SSDSE では小さすぎて並列化の恩恵がないが、 オーバーヘッドの実感を得られる教育コード。
📥 入力データ (SSDSE-B-2026 全 47 都道府県 × 12 年 2012-2023):
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 43 44 45 | import pandas as pd import time from multiprocessing import Pool df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=[1]) df.columns = ['Year','Code','Pref','Pop'] + list(df.columns[4:]) def heavy_per_pref(args): pref, group = args # 重い計算をシミュレート: 線形回帰係数 import numpy as np years = group['Year'].values pops = group['Pop'].values if len(years) < 2: return pref, None slope = np.polyfit(years, pops, 1)[0] return pref, slope groups = list(df.groupby('Pref')) # === 逐次実行 === t0 = time.time() seq = dict(heavy_per_pref(g) for g in groups) t_seq = time.time() - t0 # === 並列実行 (4 ワーカ) === # multiprocessing はワーカに関数を「名前で」渡すので、関数が import できる場所 # (.py ファイルのトップレベル)に無いと PicklingError になる。 # ノートブックや対話実行ではここで失敗するので、逐次にフォールバックする。 t0 = time.time() try: with Pool(processes=4) as pool: par = dict(pool.map(heavy_per_pref, groups)) parallel_ok = True except Exception as e: print(f'(並列実行は使えませんでした: {type(e).__name__})' '→ heavy_per_pref を別ファイルに置いて import すると動きます') par = dict(heavy_per_pref(g) for g in groups) parallel_ok = False t_par = time.time() - t0 print(f'逐次: {t_seq*1000:.1f} ms, {"並列(4)" if parallel_ok else "逐次(代用)"}: {t_par*1000:.1f} ms') print(f'Speedup = {t_seq/t_par:.2f}x') print(f'結果が一致: {seq == par}') print(f'東京の年間人口変化: {seq["東京都"]:+,.0f} 人/年') |
📤 実行例 (実測、 macOS):
💬 結果の読み方:4 ワーカが実際に動いた場合は逐次 0.9 ms に対して 555.0 ms、 並列化したのに約 600 倍遅くなった(Speedup は 0.0016 で、 小数 2 桁表示では 0.00x)。 これは SSDSE が小さすぎ、 Pool の起動 + データシリアライズ (pickle) のオーバーヘッドが本処理を上回るため。 「並列化 = 速くなる」とは限らない教訓。 1 県あたりの処理が数百 ms 以上ある場合 (例: 機械学習モデル fit) のみ、 並列化の効果が出る。
🎯 このコードでやること:Dask DataFrame で SSDSE を扱い、 「47 都道府県別の平均人口・最大年・最小年」を遅延評価で計算する。 PySpark とほぼ同じ API で、 学習教材として最適。
📥 入力データ (SSDSE-B-2026 を 8 MB ブロックで分割読込):
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 | import dask.dataframe as dd # 遅延読み込み (まだディスクは読まない) ddf = dd.read_csv( 'data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=[1], blocksize='8MB', ) # Dask のカラム名 (1 文字目を直接アサイン) ddf.columns = ['Year','Code','Pref','Pop'] + list(ddf.columns[4:]) # === 遅延グラフを組み立て === agg = ddf.groupby('Pref').agg({ 'Pop': ['mean', 'min', 'max', 'count'], 'Year': ['min', 'max'], }) # === ここで初めて実計算 (compute) === result = agg.compute() result.columns = ['Pop_mean','Pop_min','Pop_max','Pop_count','Year_min','Year_max'] print(result.sort_values('Pop_mean', ascending=False).head(5)) |
📤 実行例:
💬 結果の読み方:Dask の .compute() を呼ぶまで実計算は走らない (遅延評価)。 これにより複数の集計を組み合わせて1 回のスキャンで全部計算できる。 Spark の lazy evaluation と同じ思想。 GB〜TB 級データでも同じコードで動作し、 学習コストが極めて低い。
🎯 このコードでやること:SSDSE の 47 県データに対し、 ワーカ数 1, 2, 4, 8 で同じ集計を実行し、 Speedup をプロット。 Amdahl 法則の予測曲線と重ね合わせる。
📥 入力データ:47 都道府県別の人口リスト × 12 年分 (2012-2023)
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 43 44 45 46 47 48 49 50 51 52 | import pandas as pd import time from concurrent.futures import ProcessPoolExecutor import matplotlib.pyplot as plt import numpy as np df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=[1]) df.columns = ['Year','Code','Pref','Pop'] + list(df.columns[4:]) def heavy(args): """重い処理の代役 (12×12 行列の SVD。 実際は 1 県あたり数十マイクロ秒で終わる軽い処理)""" pref, group = args import numpy as np arr = np.array(group['Pop']) # SVD で「重い」処理を作る (本来不要) M = np.outer(arr, arr) + np.eye(len(arr)) * 1000 u, s, v = np.linalg.svd(M) return pref, s[0] groups = list(df.groupby('Pref')) speedups = [] worker_counts = [1, 2, 4, 8] t_base = None for n in worker_counts: t0 = time.time() # ProcessPoolExecutor はワーカに関数を名前で渡すので、heavy が import できる # 場所(.py ファイルのトップレベル)に無いと PicklingError になる。 # ノートブックや対話実行ではここで失敗するため、逐次にフォールバックする。 try: with ProcessPoolExecutor(max_workers=n) as ex: list(ex.map(heavy, groups)) except Exception as e: if n == worker_counts[0]: print(f'(並列実行は使えませんでした: {type(e).__name__})' '→ heavy を別ファイルに置いて import すると動きます。ここは逐次で代用') list(map(heavy, groups)) t = time.time() - t0 if t_base is None: t_base = t sp = t_base / t speedups.append(sp) print(f'workers={n}: {t*1000:.1f} ms, speedup={sp:.2f}x') # Amdahl 予測 (p=0.95 と仮定) ns = np.arange(1, 16) p = 0.95 amdahl = 1 / ((1 - p) + p / ns) plt.plot(worker_counts, speedups, 'o-', label='実測') plt.plot(ns, amdahl, '--', label=f'Amdahl (p={p})') plt.xlabel('# workers'); plt.ylabel('Speedup'); plt.legend(); plt.grid(True) plt.title('SSDSE 47 県の重い集計: 並列効率') |
📤 実行例 (実測、 macOS):
💬 結果の読み方:ワーカが実際に動くと 1 ワーカで 583.6 ms、 8 ワーカで 978.9 ms と、 増やすほど遅くなり speedup は 0.60x まで下がった。 1 県の SVD は 12×12 行列で数十マイクロ秒しかかからず、 時間のほぼ全部がプロセスの起動とデータの受け渡しなので、 p ≈ 0.95 の Amdahl 予測線(8 ワーカで約 5.9 倍)とは逆向きになる。 逐次代用の側で 1.52x まで「速く」なって見えるのは並列の効果ではなく、 1 回目に失敗したプールの起動と初回の実行コストが workers=1 だけに乗っているためである。 Amdahl の曲線どおりの伸びを見るには、 1 県あたりの計算がプロセス起動費 (数百 ms) より十分重い処理が必要。
LLM や巨大 NN は 1 GPU メモリに収まらないため、 3 種類の並列化を組み合わせる。 SSDSE 規模では不要だが、 「47 都道府県別モデル」を学習する想定で各方式の意味を見る。
| 方式 | 分割対象 | 代表ライブラリ | SSDSE での例え |
|---|---|---|---|
| Data Parallel | バッチを GPU 間で分割 | PyTorch DDP, Horovod | 47 県の人口予測モデルを 4 GPU で学習 (各 GPU が 12 県担当) |
| Model (Tensor) Parallel | 層内の行列を分割 | Megatron-LM | 1 つの巨大線形層を 4 GPU で分担 (列分割) |
| Pipeline Parallel | 層を GPU 間で分割 | GPipe, PipeDream | 前処理→特徴量抽出→予測の 3 段を 3 GPU で流れ作業 |
| FSDP (Fully Sharded) | モデル + 勾配 + Optim を全 GPU で分割 | PyTorch FSDP, ZeRO-3 (DeepSpeed) | LLM 学習の標準 |
| Expert Parallel | MoE のエキスパート群を分割 | Switch Transformer, GShard | 「東日本専門家」「西日本専門家」を別 GPU に置く |
3D 並列の例 (GPT-3 175B):
256 GPU = 8 (Pipeline) × 8 (Tensor) × 4 (Data) 1 ノード = 8 GPU (NVLink 高速接続) → Tensor Parallel ノード間 = InfiniBand → Pipeline / Data Parallel バッチ = 1 ノードに data 分散、 ノード内で tensor 分散
SSDSE 視点: 47 都道府県の人口予測モデル (LSTM, 100 万パラメタ) ならData Parallel だけで十分。 「都道府県をバッチ次元として 8 GPU で分散」と理解すれば、 巨大 LLM の Data Parallel と本質は同じ。
SSDSE-B-2026 は年度集計データ (バッチ) だが、 もし「47 都道府県の住民票異動を秒単位で受信」する仮想システムを設計するなら、 ストリーム処理パラダイムが必要になる。
| 処理方式 | 遅延 | スループット | 代表 OSS |
|---|---|---|---|
| バッチ | 時間〜日 | 最大 | Hadoop, Spark batch, BigQuery |
| マイクロバッチ | 秒 | 大 | Spark Streaming |
| レコード単位ストリーム | ミリ秒 | 中 | Flink, Storm, Kafka Streams |
| CEP (Complex Event) | サブミリ秒 | 中 | Esper, Flink CEP |
ストリーム処理の 3 大課題:
Lambda アーキテクチャ vs Kappa アーキテクチャ:
| アーキテクチャ | 構成 | 利点 | 欠点 |
|---|---|---|---|
| Lambda | バッチ層 + 速度層の 2 系統 | 正確性 + 低遅延 | 2 重実装が大変 |
| Kappa | ストリーム 1 本 | シンプル | 過去再処理がコスト高 |
SSDSE 規模なら自前 PC で十分だが、 「47 都道府県別の住民全データ (約 1.2 億行 × 数百カラム)」を扱うなら、 マネージドサービスを選ぶのが現実的。 主要サービスを比較する。
| 用途 | AWS | GCP | Azure |
|---|---|---|---|
| DWH | Redshift | BigQuery | Synapse |
| マネージド Spark | EMR / Glue | Dataproc / Dataflow | HDInsight / Synapse Spark |
| ストリーム | Kinesis / MSK | Pub/Sub / Dataflow | Event Hubs / Stream Analytics |
| 分散 NoSQL | DynamoDB | Bigtable / Firestore | Cosmos DB |
| 分散 ML 学習 | SageMaker | Vertex AI | ML Studio |
| サーバレス分散 | Lambda + Fargate | Cloud Run / Functions | Functions / Container Apps |
| 分散ストレージ | S3 | GCS | Blob Storage |
| 分散コンセンサス | Step Functions | Workflows | Logic Apps |
料金感 (2026 年想定、 例): SSDSE-B-2026 (約 0.36 MB) を BigQuery にロード→47 都道府県集計→ Looker Studio で可視化のフローは月額 < $0.01。 一方、 全国民住民票 (約 100 GB) なら月額数千ドル規模。
multiprocessing の意味はない。Dask/Spark は GB 級から効く。大きな集計タスクを「分けて・同時に・まとめる」を体感するシミュレータ。 架空の生成データ (固定シードの擬似乱数 4,700 個、 値域 100〜999) の総合計を、 1〜8 台のワーカで分担して計算する。 各ワーカの部分和 (Map)、 統合 (Reduce)、 仮想処理時間はすべて実際に JavaScript で正確に計算しており、 何台に分けても最終合計は必ず同じ値になる。
※ データは架空 (固定シード擬似乱数)。 時間モデルも仮想: 分割などの逐次オーバーヘッド 100 ms + 1 個あたり処理 0.5 ms + ワーカ 1 台あたり通信・統合 35 ms。 部分和・総合計はチャンクの実値から正確に計算。 グラフは故障なしの理論時間、 ● が現在のワーカ数、 ✕ が直前の実測 (故障ありは遅くなる)。
4,700 枚の答案の採点を 1 人でやれば 4,700 枚分の時間がかかるが、 8 人で分ければ 1 人あたり約 588 枚。 これが Map (分割して並列に部分処理)。 ただし最後に 8 人分の小計を集めて合算する人が要る — これが Reduce (統合)。 「配る手間」と「集める手間」は人数を増やしても消えない。 むしろ人数に比例して増える。 だから台数を 2 倍にしても時間は半分にならない。 上のシミュレータで 1 台→2 台は大きく縮むのに、 7 台→8 台がほとんど縮まないのはこのためである。
groupby / join ではワーカ間を入力データの 0.5〜2 倍のバイトが飛び交い、 ネットワークが律速になる。 事前集約 (combiner) で送る量を小さくするのが定石。このシミュレータの「分割 → 並列部分集計 → 統合」は Hadoop / MapReduce (Google, 2004) の骨格そのもの。 Apache Spark は中間結果をメモリに保持して同じパラダイムを高速化し、 系統情報 (lineage) により故障時は失われたパーティションだけ再計算する — 故障トグルで見た「担当分だけやり直す」の洗練版である。 頭打ちの理論がアムダールの法則 $S(N) \le 1/((1-p) + p/N)$: 並列化できない割合 $(1-p)$ が上限を決める ($p=0.95$ なら何台並べても最大 20 倍)。 データ量も一緒に増やす現実的な見方が Gustafson の法則で、 ビッグデータ 時代に分散処理が有効な理由はこちらで説明される。
データ分散処理 ★
├─ パラダイム
│ ├─ MapReduce (Hadoop)
│ ├─ DAG実行 (Spark, Dask)
│ ├─ ストリーム (Flink, Beam)
│ └─ Actor (Ray, Akka)
├─ ストレージ
│ ├─ 分散FS (HDFS / S3 / GCS)
│ ├─ 列指向 (Parquet / ORC)
│ └─ メッセージング (Kafka)
├─ 性能理論
│ ├─ Amdahl (上限あり)
│ └─ Gustafson (問題拡大)
└─ 整合性
├─ CAP 定理
├─ コンセンサス (Raft / Paxos)
└─ 結果整合性 (eventual)
分散データ処理はエンジニアリングと MLOps の中核技術。
データ分割 → 並列処理 → 中間結果統合 → 最終集計の流れで、 partition pruning とシャッフル最小化がパフォーマンスを決める。
分散データ処理は規模・モード・基盤の三段で技術選定する。
小〜中規模なら pandas、 大規模バッチなら Spark、 超低レイテンシなら Flink、 と「規模とリアルタイム性」で選ぶ。
分散処理の核心は、 1 台の性能を上げる (スケールアップ) のではなく、 普通のマシンを台数で並べる (スケールアウト) という発想の転換にある。 CPU を速くしたり RAM を増やしたりする道は物理法則と価格で頭打ちになるが、 「安いマシンを 100 台つなぐ」道はほぼ線形にコストを積めば容量が伸びる。 データと計算を複数マシンに撒き、 同じ処理を同時に走らせて (並列)、 最後に結果を寄せ集める — これだけで単一マシンの限界を超えられる。 MapReduce と Apache Spark はこの「撒く → 各自で処理 → 集める」を汎用化した骨格である。
| 観点 | スケールアップ (垂直) | スケールアウト (水平・分散) |
|---|---|---|
| 増やすもの | 1 台の CPU/RAM/ディスク | マシンの台数 (ノード) |
| 上限 | 物理・価格で急に頭打ち | 台数を足せば伸びる (準線形) |
| 故障の影響 | その 1 台が単一障害点 | 1 台落ちても全体は継続 (要設計) |
| 向くデータ規模 | RAM に収まるうち | RAM・単一ディスクを超える規模 |
大事なのは「分散すれば速い」ではないこと。 撒く・集めるには通信という追加コストが必ず生じる。 だから分散処理は「1 台では終わらない、 または現実的な時間に収まらない」ときにだけ効く道具である。 SSDSE-B-2026 (実測 564 行 × 112 列、 ファイル約 351 KB) のように 単一マシンの RAM に楽々収まる規模では、 pandas 一発が最速で、 分散はむしろ遅くなる。 本ページが 47 都道府県を「47 ノード」に見立てるのは、 発想をローカルで安全に体験するためのメタファーだと割り切ってほしい。
「台数を増やせば増やすほど速い」は幻想である。 分散処理には単一マシンには存在しない固有のコストと不確実性がまとわりつく。 ここでは特に効きの大きい 7 点を、 具体的な現れ方とともに整理する。
groupby / join / sort は同じキーのデータを 1 か所に集めるため、 ワーカ間を大量のバイトが飛ぶ (シャッフル)。 これはネットワーク帯域が律速で、 台数を増やしても縮まないどころか増えることもある。 事前集約 (combiner)、 broadcast join、 partition キーの事前整列で「送る量」を減らすのが定石。スキューの効き目を実データで確かめる。 SSDSE-B-2026 の 2023 年 47 都道府県について、 各県の総人口 (A1101) を「その県シャードの処理量」の代理とみなす。 実測では最大の東京都 14,086,000 に対し最小の鳥取県 537,000 で、 約 26.2 倍の開きがある (これがスキューの正体)。 これを 4 ワーカに割り当てる 2 方式を比べる。 いずれの負荷値も上記 CSV から実際に計算した実測値である。
| 分割方式 | 各ワーカの負荷 (総人口の和) | 最大/最小 |
|---|---|---|
| 素朴な連番 4 分割 (コード順にそのまま切る) | 33,622,000 / 45,791,000 / 28,028,000 / 16,912,000 | 2.71 |
| 貪欲な均等割り (負荷の大きい県から最軽量ワーカへ / LPT法) | 31,223,000 / 31,142,000 / 31,026,000 / 30,962,000 | 1.01 |
素朴分割では最も重いワーカが最も軽いワーカの 2.71 倍働くため、 全体の完了は重いワーカに引きずられる。 一方、 大きい県から順に「今いちばん軽いワーカ」へ配る貪欲法 (LPT: Longest Processing Time first) を使うと 1.01 倍までならせ、 ほぼ理想の負荷分散になる。 実務ではキーを hash して均等に振る、 重いキーを分割する (salting)、 動的リバランスといった手法で同じ効果を狙う。 分散の速さは「いちばん遅いワーカ」で決まるという原則が、 この実測からそのまま読み取れる。
※ 上表の負荷値は SSDSE-B-2026 (2023 年 47 都道府県) の総人口を実際に集計した実測。 「総人口=処理量」という対応づけ自体は説明用のモデル (スキューを可視化するための仮定) であり、 実データの数値そのものは捏造ではない。 貪欲均等割りの割り当て手順は架空の教材例。
基本パラダイム (map → shuffle → reduce) の先には、 ストレージ・並列化の軸・整合性・最適化という広い世界がある。 ここでは代表トピックを俯瞰し、 深掘り先の関連ページへ橋渡しする。
大規模データはまず「どこに置くか」で決まる。 HDFS は巨大ファイルをブロックに割り、 複数ノードにレプリケーション (通常 3 重) して保存する。 計算をデータの近くへ送る「データ局所性」を前提に設計され、 Hadoop エコシステムの土台になった。 クラウドでは S3/GCS などのオブジェクトストレージが同じ役割を担う。 列指向フォーマット (Parquet / ORC。 本ページ内リンクなし) は必要列だけ読めるため、 CSV より桁違いに速い。
分散機械学習では並列化の軸が 2 つある。 データ並列は「同じモデルを各ワーカに複製し、 データを分割して勾配を平均する」方式で、 大量データに向く (Horovod, FSDP, DDP)。 モデル並列は「巨大すぎて 1 台に載らないモデル自体を層やテンソルで切って複数ワーカに配る」方式で、 LLM のような巨大モデルに必須 (テンソル並列・パイプライン並列)。 実際の大規模学習は両者を組み合わせる。 運用面は AI クラウド や MLOps が扱う。
CAP 定理で可用性 (AP) を選ぶと、 一時的に古い値が返るがいずれ全レプリカが同じ値に収束するという「結果整合性」を受け入れる。 SNS のいいね数のように「多少ズレても後で揃えばよい」用途に向く。 逆に会計や在庫のように即時の強整合が要る用途は CP を選ぶ。 「どちらを選ぶか」は 分散 DB 選定の中心的論点である。
分散処理の性能はほぼシャッフル量で決まるため、 最適化の主戦場もそこにある。 パーティショニングで join/groupby のキーを事前に揃えておけば、 再配置なしで処理でき (co-partitioned join)、 小さい参照表は全ワーカに配る (broadcast join) とシャッフルを消せる。 事前集約 (combiner)、 partition pruning (不要ブロックを読まない)、 適切な partition 数 (1 partition = 数百 MB〜数 GB が経験則) の設計が効く。 Spark の Catalyst/Tungsten、 ビッグデータ基盤のクエリ最適化はこれらを自動化する方向に進化してきた。
※ Parquet / ORC / Kafka などは本用語集に個別ページが無いためテキスト表記に留めた。