「Spark」を取り巻く中核キーワード群です。 検索やインデックス作成で参照する際の手がかりにしてください。 各キーワードは関連する概念・手法・道具立てを含み、 文献検索や学習計画の起点になります。
🍰 まずはやさしく
大量のデータを処理する高速な道具です。
たくさんの機械で分担して計算するために使います。
スマホの膨大な利用データを分析する時に便利です。
この章ではSparkの結論をまとめます。
最も忙しい読者のために、 まず結論だけまとめます。 詳細は以下のセクションへ:
🍰 まずはやさしく
データの量が多くて困った時の救世主です。
1台のパソコンで処理できないデータを扱うために使います。
数千万行もある名簿を集計するような場面で出番があります。
この章ではSparkをいつ使うかを説明します。
「数千万行の CSV を集計したい」 「pandas だとメモリに乗らない」 — そんなとき Spark の出番。 クラスタ(複数台のマシン)にデータを分散し、 並列で処理します。
このページの読み方:まず 30秒結論 と 直感 を読み、 必要に応じて 数式 や 計算例、 落とし穴 に進んでください。
🍰 まずはやさしく
大勢で分担して作業するチームのようなものです。
計算時間を短くするために使います。
1万冊の本を100人で分担して読むイメージです。
この章ではSparkが動く仕組みを図で解説します。
1 万冊の蔵書を読みたいとき:
Spark の役割は「誰がどの本を読むか」 「読んだ結果をどう集めるか」を取り仕切ること。 ユーザーは df.groupBy(...).count() と書くだけで、 裏では 100 台のマシンが並列で動きます。
Spark の3つの中核概念 (クラスタ構成 / RDD 系統グラフ / Stage 実行) を inline SVG で視覚化する。 これらの図は本文の言葉だけでは把握しづらい「分散処理の流れ」「変換と行動の境界」「DAG スケジューラの仕事」を一枚に集約したものである。
Driver プログラムが SparkContext を作成し、 Cluster Manager (YARN / Kubernetes / Standalone) を介して複数の Worker 上に Executor を立ち上げる。 タスクは Driver → Executor へ送られ、 結果は逆方向に集約される。
Driver は「司令塔」、 Cluster Manager は「リソース調整役」、 Executor は「実働部隊」。 タスクが何百あっても、 各 Executor のスレッド数だけ並列に処理される。
RDD は「親 RDD と変換ルールの記録」。 map / filter / flatMap は遅延評価で、 collect / count / save 等の action が来た時点で初めて DAG 全体が実行される。 これにより最適化と耐障害性 (再計算可能) を両立する。
reduceByKey など widedependency (shuffle) を挟む変換が Stage の境界となる。 系統が分かれば、 Executor が落ちても親 RDD から再構築できる。
action が呼ばれると Job が作られ、 shuffle 境界で Stage に分割、 各 Stage は partition 数だけ Task に分かれて Executor で並列実行される。 この階層が Spark UI の表示構造そのものである。
Stage 1 (例: map+filter) は shuffle なしで完結し、 Stage 2 (例: reduceByKey 後段) は前段の shuffle ファイルを読み込んでから走る。 partition 数 = Task 数なので、 partition を増やせば並列度が上がるが、 タスク起動コストとのバランスが必要。
この3枚を脳内に置けば、 Spark UI の「Jobs / Stages / Tasks」タブを開いたときに何を見ているかが瞬時にわかる。 さらに学習を進めるなら Hadoop・分散処理・クラウドサービス ページへ。
Apache Spark の中核を本当に掴めたかを 6 問で確認する。 すべて 1 分以内で答えられる粒度。 「自分で」答えを口に出して言える状態を目指したい。
map, filter, select)は 遅延評価で DAG を構築するだけ。 Action(count, collect, write)が呼ばれて初めて実行される。
groupByKey, reduceByKey, join, repartition, distinct などキーの再分配が必要な操作。 shuffle はネットワーク・ディスク I/O が重く、 性能のボトルネックになりやすい。
これらの問題に詰まった項目があれば、 「DAG 構造」「shuffle の仕組み」「Catalyst Optimizer」の該当セクションへ戻ろう。 Spark は「いつ使うべきか」の判断こそが熟練度のバロメータである。
Spark の心臓部である 遅延評価(lazy evaluation) と パーティション並列 を、 12 個の数値からなる 架空データ(説明用・決定的)で体感します。 map/filter などの transformation を積んでもすぐには計算されず、 計算グラフ(DAG)が伸びるだけ。 count/collect などの action を呼んだ瞬間に一斉実行される — この境界を目で見て掴んでください。 Hadoop のワードカウント(Map→Shuffle→Reduce)とは別の側面に絞っています。
機械学習やグラフ計算のように 同じデータを何度も舐める反復処理こそ Spark の主戦場。 Hadoop MapReduce は反復のたびに中間結果を HDFS(ディスク)へ書き戻すのに対し、 Spark は cache() でメモリ保持し、 読み書きを省く。 反復回数を変えて I/O 差を見てください。
rdd.map(f) を書いても f は実行されず、 count() を呼ぶまで例外もログも出ない。 バグの発火位置が action 行までズレる。 ②action を複数回呼ぶと DAG が毎回再計算される — 同じ RDD を使い回すなら cache()/persist() しないと、 collect と count で 2 回フル実行される。 ③shuffle は並列で割れない — groupByKey/join はパーティション間の全対全通信を伴い、 上のスライダーで並列度を上げても縮まらない部分が残る。 ④メモリ不足 — 保持しきれないと disk へ溢れ(spill)、 Hadoop 同様に遅くなる。 インメモリは万能ではない。🍰 まずはやさしく
データから答えを導き出すための仕組みです。
正しく数値を計算したり分類したりするために使います。
テストの点数からクラスの傾向を出すような処理です。
この章ではSparkの定義を数式で説明します。
spark の定義や代表的な数式を以下に示す。 数式の各記号の意味は次節で言葉に翻訳する。
spark は文脈に応じて複数の定式化があるが、 教育目的では最も基本的な形を抑えることが重要。 具体的な値での計算例は後続セクションを参照。
$$\text{spark}: f(\mathbf{X}, \boldsymbol{\theta}) \to y$$
記号の対応はこうです。 Spark の実行単位は上位から Job → Stage → Task の 3 層で、 $\text{Stage}$ の切れ目はシャッフル(ノード間のデータ移動)が必要になる箇所に対応します。 $\text{Task}_{1..k}$ の $k$ はパーティション数で、 1 パーティション=1 タスク=1 コアで処理される単位です。 スループットの式 $\text{Data Size} / (\text{Cluster Cores} \times \text{Time per partition})$ が示すのは、 コアを増やしても、パーティション数がコア数より少なければ余ったコアは遊ぶということ。 経験則としてパーティション数はコア数の 2〜4 倍にします。 SSDSE の 564 行のような小さなデータでは、 パーティションを分ける費用のほうが計算より高くつきます。
map, filter, groupBy 等。 遅延評価(呼んでもまだ実行されない)。count, show, collect 等。 これを呼んだ瞬間に実行が走る。「Apache Spark」は 大規模分散データ処理の業界標準。インメモリ計算で Hadoop MapReduce より 10〜100 倍高速。 です。 ここでは定義式の各記号、 直感的意味、 SSDSE-B-2026 への当てはめを段階的に解きほぐします。
| モジュール | 機能 | SSDSE-B-2026 対応 |
|---|---|---|
| Spark Core / SQL | 分散データフレーム | 47 都道府県 ETL |
| Spark Streaming | リアルタイム処理 | 毎時人口移動データ |
| MLlib | 分散機械学習 | 都道府県分類モデル |
| GraphX | グラフ処理 | 都道府県間関係 |
SparkSession を取得 → DataFrame 読み込み → transformation (select, filter, groupBy) → action (show, count, write) の 4 ステップが基本。 SQL クエリも spark.sql("SELECT ...") で実行可能。
SSDSE-B-2026 を PySpark で都道府県別集計。 入力は 47 都道府県 × 12 年分の SSDSE-B-2026 行、 出力は Spark DataFrame で groupBy・agg・SQL クエリ。 47 都道府県 × 複数年の実データで具体計算を実施します。
| 業界 | 事例 | 役割 | SSDSE-B-2026 との対比 |
|---|---|---|---|
| Netflix | 視聴ログ 1PB/日のリアルタイム集計 | 推薦エンジン学習基盤 | 47 都道府県データを 47 パーティションで並列処理 |
| Uber | 数十億のライドデータ ETL | Spark SQL + MLlib | 都道府県別輸送統計に応用可 |
| Facebook/Meta | ソーシャルグラフ解析 GraphX | 影響範囲解析 | 都道府県人口移動グラフに応用 |
| Yahoo Japan | 広告ログ集計とリアルタイム最適化 | Structured Streaming | SSDSE-B の年次集計バッチに類似 |
| 画像レコメンドのバッチ学習 | MLlib + Tensorflow | 出生率予測などの ML 基盤 | |
| Goldman Sachs | リスク計算ジョブ並列化 | DataFrame API | 金融指標の時系列分析 |
| 手法 | 定義 | 特徴 | 用途 |
|---|---|---|---|
| Apache Spark | インメモリ分散処理 | JVM、 Python (PySpark) API | 汎用大規模 ETL/ML |
| Hadoop MapReduce | ディスクベースの分散処理 | 古典、10倍以上遅い | 歴史的基盤 |
| Apache Flink | ストリーム優先 | 低遅延・状態管理 | リアルタイム解析 |
| Dask | Python ネイティブ並列計算 | pandas/numpy 互換 | 中小規模 Python 分析 |
| Ray | AI/強化学習向け分散 | Python・低レイテンシ | ML 推論サーバ |
| Presto/Trino | 対話型 SQL クエリ | クエリ専用、ETL に弱い | BI ダッシュボード |
| BigQuery | Google マネージドサーバレス | サーバ管理不要 | クラウド分析 |
collect() で集めるwrite でファイル出力するか take(n) で先頭サンプルだけ取る。spark.sql.shuffle.partitions を core 数の 2-4 倍に。broadcast() を明示。pyspark.sql.functions のビルトインを使う。 必要なら Pandas UDF (vectorized) で 10-100倍高速化。1 2 3 4 5 6 7 8 | from pyspark.sql import SparkSession spark = SparkSession.builder.appName('SSDSE').getOrCreate() import pandas as pd # SSDSE-B-2026.csv は 1 行目が項目コード・2 行目が日本語見出しの 2 段ヘッダ。 # Spark の CSV 読み込みでは 1 行目を飛ばせないので、pandas で読んでから Spark に渡す pdf = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=1) df = spark.createDataFrame(pdf) df.filter(df['年度']==2023).select('都道府県','総人口').show(5) |
df.createOrReplaceTempView('t')\nspark.sql(\"SELECT 都道府県, SUM(出生数) AS s FROM t GROUP BY 都道府県 ORDER BY s DESC\").show(10)
partitionBy しておく。 (2) 集計関数を reduceByKey 系で局所集約。 (3) salting で偏り解消。 (4) Pre-aggregation でデータ量削減。
1 2 3 4 5 6 | from pyspark.ml.feature import VectorAssembler from pyspark.ml.regression import LinearRegression va = VectorAssembler(inputCols=['総人口'], outputCol='features') d = va.transform(df).select('features','出生数') lr = LinearRegression(labelCol='出生数').fit(d) print(lr.coefficients, lr.intercept) |
PySpark で都道府県データを集計:
| 操作 | pandas | PySpark |
|---|---|---|
| 読込 | pd.read_csv | spark.read.csv |
| フィルタ | df[df.x>0] | df.filter(df.x>0) |
| 集計 | df.groupby().mean() | df.groupBy().mean() |
| 実行タイミング | 即時 | action まで遅延 |
合成 1TB データを N ノードで処理する所要時間を計算する。
1 2 3 4 5 6 7 | import numpy as np data_mb = 1_000_000 speed = 100 N = np.array([1, 10, 100]) serial = data_mb / speed parallel = serial / N print(f"並列時間: {parallel/60} 分") |
💬 手計算 (Step 2) と Python 出力が完全一致。
最小再現コード。 SSDSE-B のような実データを前提に、 4〜8 行で動く例です:
1 2 3 4 5 6 7 8 9 | from pyspark.sql import SparkSession spark = SparkSession.builder.appName('demo').getOrCreate() import pandas as pd # SSDSE-B-2026.csv は 1 行目が項目コード・2 行目が日本語見出しの 2 段ヘッダ。 # Spark の CSV 読み込みでは 1 行目を飛ばせないので、pandas で読んでから Spark に渡す pdf = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=1) df = spark.createDataFrame(pdf) df.groupBy('年度').count().orderBy('年度').show() # action で実行される spark.stop() |
補足:ライブラリのバージョンや前処理状態によって出力は変わります。 自分の環境で動かすときは pip list でバージョンを確認し、 入力 CSV のパス・列名を実態に合わせてください。
🎯 このコードでやること:47 都道府県 × 12 年分のデータを Spark DataFrame として読み込み、 年度別に総人口を集計する。
📥 入力データ:SSDSE-B-2026.csv (cp932, 564 行)、 年度・都道府県・総人口を含む。
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 | from pyspark.sql import SparkSession from pyspark.sql.functions import sum as ssum, avg spark = (SparkSession.builder .appName('SSDSE-Analysis') .config('spark.sql.shuffle.partitions', '4') .getOrCreate()) import pandas as pd # SSDSE-B-2026.csv は 1 行目が項目コード・2 行目が日本語見出しの 2 段ヘッダ。 # Spark の CSV 読み込みでは 1 行目を飛ばせないので、pandas で読んでから Spark に渡す pdf = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=1) df = spark.createDataFrame(pdf) # 文字列→数値 df = df.withColumn('総人口', df['総人口'].cast('long')) df = df.withColumn('年度', df['年度'].cast('int')) print("行数:", df.count()) (df.groupBy('年度') .agg(ssum('総人口').alias('全国総人口'), avg('総人口').alias('県平均')) .orderBy('年度') .show()) |
📤 実行結果(Spark では未実行:Spark 実行環境が無いため、 同じ集計を pandas で実測し、 show() の表の形に整えた値):
💬 結果の読み方:Spark DataFrame で 564 行を読み込み、 12 年分の全国総人口推移を 1 ジョブで集計。 2012→2023 で約 324 万人減少 (年率 約-29.4 万人)。 同じことを pandas でやっても可能だが、 これが TB 級になると Spark でないと不可能。
🎯 このコードでやること:Spark SQL を使い 2023 年の都道府県別合計特殊出生率を降順表示。
📥 入力データ:df: SSDSE-B-2026 全データ。 「合計特殊出生率」列を含む。
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 spark = SparkSession.builder.getOrCreate() import pandas as pd # SSDSE-B-2026.csv は 1 行目が項目コード・2 行目が日本語見出しの 2 段ヘッダ。 # Spark の CSV 読み込みでは 1 行目を飛ばせないので、pandas で読んでから Spark に渡す pdf = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=1) df = spark.createDataFrame(pdf) df.createOrReplaceTempView('ssdse') result = spark.sql(""" SELECT 都道府県, CAST(合計特殊出生率 AS DOUBLE) AS tfr FROM ssdse WHERE 年度 = 2023 ORDER BY tfr DESC LIMIT 10 """) result.show() |
📤 実行結果(Spark では未実行:Spark 実行環境が無いため、 同じ集計を pandas で実測し、 show() の表の形に整えた値):
💬 結果の読み方:Spark SQL で 2023 年の出生率 TOP 10 を抽出。 沖縄が突出して高く 1.60、 西日本・九州が高い傾向。 同じクエリが 47 都道府県でも 47 億行でも書き換え不要 — これが Spark の力。
🎯 このコードでやること:47 都道府県の 2023 年データで、 総人口から出生数を予測する線形回帰モデルを Spark MLlib で学習。
📥 入力データ:X = 総人口, y = 出生数, 47 行。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 | from pyspark.ml.feature import VectorAssembler from pyspark.ml.regression import LinearRegression from pyspark.sql.functions import col d = (df.filter(col('年度')==2023) .withColumn('総人口', col('総人口').cast('double')) .withColumn('出生数', col('出生数').cast('double')) .select('総人口', '出生数')) va = VectorAssembler(inputCols=['総人口'], outputCol='features') train = va.transform(d).select('features', '出生数') lr = LinearRegression(labelCol='出生数', featuresCol='features') model = lr.fit(train) print(f"係数: {model.coefficients[0]:.6f}") print(f"切片: {model.intercept:.2f}") print(f"R^2 : {model.summary.r2:.4f}") |
📤 実行結果(Spark では未実行:Spark 実行環境が無いため、 同じ集計を pandas で実測し、 show() の表の形に整えた値):
💬 結果の読み方:総人口 1000 万人増 → 出生数 約 6.10 万人増 (出生率約 6.10‰)。 R²=0.991 と極めて高い決定係数。 sklearn と同じ結果が Spark でも得られる。 数十億行の住民データでも同じコードで動く点が Spark の意義。
🎯 このコードでやること:SSDSE-B-2026 を年度別バッチとして処理し、 各年度の全国合計・最大・最小を出力 (ストリーミング処理の擬似デモ)。
📥 入力データ:df: 全データ、 年度別にループ。
1 2 3 4 5 6 7 8 9 10 | from pyspark.sql.functions import sum as ssum, max as smax, min as smin, col for yr in sorted([r['年度'] for r in df.select('年度').distinct().collect()]): sub = df.filter(col('年度')==yr) agg = sub.agg( ssum('総人口').alias('合計'), smax('総人口').alias('最大'), smin('総人口').alias('最小') ).first() print(f"{yr}: 合計={agg['合計']:>10,}, 最大={agg['最大']:>9,}, 最小={agg['最小']:>7,}") |
📤 実行結果(Spark では未実行:Spark 実行環境が無いため、 同じ集計を pandas で実測し、 show() の表の形に整えた値):
💬 結果の読み方:12 年分を 1 ループで集計。 最大は常に東京 (1300万→1400万)、 最小は鳥取 (58万→54万)。 ストリーミングではこのループが Kafka など実時間ストリームになる。 Spark Structured Streaming なら全く同じ DataFrame API で書ける。
local[*] モードで全 CPU コアを使う。 開発・テストには便利。 ただし起動オーバーヘッドあり。unpersist() 明示。 または OOM 時に Spark が自動 evict。 LRU が基本。cast('long') 明示。df.select([sum(c) for c in df.columns]) または df.agg(*[sum(c) for c in cols])window 関数で時間ウィンドウ集計。 Structured Streaming で event-time、 watermark もサポート。spark.executor.memory を増やす、 partition 数を増やす、 cache を減らす、 broadcast 化、 skew 対応。変換は DAG に蓄積されるだけ。 action (count/show/write) で初めて実行。 不要計算を省くため。
ネットワーク I/O が最大コスト。 100GB の shuffle ≒ 1 ノード I/O 10 分。 broadcast join で回避。
WHERE 句を JOIN の前に押し下げる (predicate pushdown)。 列選択も pushdown。
DataFrame は Catalyst + Tungsten で型最適化&コード生成。 同じ処理で 2-10 倍速い。
Lineage (来歴) を保持。 ノード故障時にパーティションを再計算。 チェックポイントで lineage 短縮。
SSDSE-B-2026 を使った段階的ハンズオン。 初学者→中級→上級と順に深めるシナリオ構成です。
SparkSession.builder.appName('SSDSE').getOrCreate()。 local モードでも動く。
spark.read.option('header',True).option('encoding','cp932').csv('data/raw/SSDSE-B-2026.csv')。
df.withColumn('総人口', col('総人口').cast('long'))。 CSV は全列 string なので明示。
df.filter(col('年度')==2023) で 2023 年に絞る。
df.groupBy('年度').agg(sum('総人口').alias('全国'))。
spark.sql('SELECT * FROM t WHERE 年度=2023 ORDER BY 総人口 DESC LIMIT 5')。
Window.partitionBy('年度').orderBy(col('総人口').desc()) で年度別ランキング。
VectorAssembler + LinearRegression で「総人口 → 出生数」を学習。
df.write.partitionBy('年度').parquet('out/')。
http://localhost:4040 で Job/Stage/Task の所要時間とメモリを可視化。
実務で頻用するコード・概念・公式を 1 ページにまとめた早見表。 印刷して机に貼っておくと便利です。
| 領域 | 項目 | 説明 |
|---|---|---|
| 起動 | SparkSession.builder.appName('app').getOrCreate() | SparkSession 取得 |
| 読み込み | spark.read.csv('path', header=True) | CSV 読み込み |
| 表示 | df.show(5) | 先頭 5 行表示 |
| 件数 | df.count() | 全行数 |
| カラム | df.columns | カラム名リスト |
| スキーマ | df.printSchema() | 型情報 |
| フィルタ | df.filter(df['年度']==2023) | 条件抽出 |
| 選択 | df.select('都道府県', '総人口') | 列選択 |
| 集計 | df.groupBy('年度').sum('総人口') | グループ集計 |
| SQL | df.createOrReplaceTempView('t') | ビュー作成 → SQL クエリ可 |
| ソート | df.orderBy('総人口', ascending=False) | ソート |
| 結合 | df1.join(df2, 'key') | DataFrame 結合 |
| Broadcast 結合 | df.join(broadcast(small), 'key') | 小テーブル shuffle 回避 |
| キャッシュ | df.cache() | メモリキャッシュ |
| パーティション | df.repartition(8) | パーティション数変更 |
| 書き込み | df.write.parquet('out/') | Parquet 出力 |
| 到 pandas | df.toPandas() | pandas DataFrame 化 |
| UDF | from pyspark.sql.functions import udf | ユーザ定義関数 |
| Pandas UDF | @pandas_udf('long') | vectorized UDF (高速) |
| MLlib | from pyspark.ml.feature import VectorAssembler | ML 用ベクタ化 |
本編のコードに加え、 さらに 4 種の発展的コード例を 4 要素 (🎯/📥/📤/💬) 付きで提示。 段階的に「読む→動かす→改造する」を体験できます。
🎯 このコードでやること:各年度内で都道府県の総人口ランキングを Window 関数で計算。
📥 入力データ:SSDSE-B-2026 の関連カラム。
1 2 3 4 5 6 7 8 | from pyspark.sql.window import Window from pyspark.sql.functions import row_number, col w = Window.partitionBy('年度').orderBy(col('総人口').cast('long').desc()) df.withColumn('rank', row_number().over(w)) \ .filter(col('rank') <= 3) \ .select('年度','都道府県','総人口','rank') \ .show(15) |
📤 実行結果(Spark では未実行:Spark 実行環境が無いため、 同じ集計を pandas で実測し、 show() の表の形に整えた値):
💬 結果の読み方:Window 関数で年度別 TOP 3 を抽出。 全年度で東京・神奈川・大阪が不動。 SQL の RANK と同じ機能を PySpark で実現。
🎯 このコードでやること:Pandas UDF (vectorized) で人口規模カテゴリを高速付与。
📥 入力データ:SSDSE-B-2026 の関連カラム。
1 2 3 4 5 6 7 8 9 10 | from pyspark.sql.functions import pandas_udf import pandas as pd @pandas_udf('string') def categorize(pop: pd.Series) -> pd.Series: # pd.cut はカテゴリ型を返すので、'string' 型の UDF では文字列に直して返す return pd.cut(pop.astype(int), bins=[0, 1000000, 3000000, 10000000, 1e10], labels=['小','中','大','超大']).astype(str) df = df.withColumn('カテゴリ', categorize(col('総人口'))) df.groupBy('カテゴリ').count().show() |
📤 実行結果(Spark では未実行:Spark 実行環境が無いため、 同じ集計を pandas で実測し、 show() の表の形に整えた値):
💬 結果の読み方:合計 115+329+108+12=564 で、 年度で絞っていないので 12 年度分の県×年を数えている。 「超大」の 12 件は東京都の 12 年度分だけで、 1,000 万人を超えるのは東京都しかない。 Pandas UDF は列をまとめて pandas に渡すので 1 行ずつ呼ぶ通常の UDF より速いが、 pd.cut の戻り値はカテゴリ型なので、 型を 'string' と宣言したら文字列に直してから返す。
🎯 このコードでやること:VectorAssembler → StandardScaler → LogisticRegression をパイプライン化。
📥 入力データ:SSDSE-B-2026 の関連カラム。
1 2 3 4 5 6 7 8 9 10 | from pyspark.ml import Pipeline from pyspark.ml.feature import VectorAssembler, StandardScaler from pyspark.ml.classification import LogisticRegression va = VectorAssembler(inputCols=['総人口','出生数'], outputCol='raw') sc = StandardScaler(inputCol='raw', outputCol='features') lr = LogisticRegression(labelCol='label', featuresCol='features') pipe = Pipeline(stages=[va, sc, lr]).fit(train_df) preds = pipe.transform(test_df) preds.select('都道府県','probability','prediction').show(5) |
📤 実行結果(train_df・test_df を用意していない擬似コードなので、 表示の形を示す例で実測値ではない):
💬 結果の読み方:Pipeline で前処理 + 学習 + 推論を一気通貫化。 production デプロイでもこのパイプラインをそのまま使える。 sklearn とほぼ同じ API。
🎯 このコードでやること:SSDSE-B のファイルを 1 行ずつ追加する疑似ストリーミングを構成。
📥 入力データ: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 | from pyspark.sql import SparkSession spark = SparkSession.builder.getOrCreate() import pandas as pd # SSDSE-B-2026.csv は 1 行目が項目コード・2 行目が日本語見出しの 2 段ヘッダ。 # Spark の CSV 読み込みでは 1 行目を飛ばせないので、pandas で読んでから Spark に渡す pdf = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=1) df = spark.createDataFrame(pdf) # 列名と型(スキーマ)をここから借りる # Structured Streaming のテンプレート (実環境で動作) stream = (spark.readStream .option('header', True) .option('encoding', 'cp932') .schema(df.schema) .csv('stream_dir/')) query = (stream.groupBy('年度') .count() .writeStream .outputMode('complete') .format('console') .start()) |
📤 実行例(イメージ・未実測:Spark 実行環境が無く、 stream_dir/ にファイルを置いたときに出る表示の形だけを示す):
💬 結果の読み方:start() は監視を始めるだけで、 その時点では何も表示されない。 stream_dir/ に CSV が置かれるたびに console へ「Batch: 0」「Batch: 1」… と年度別の件数表が出力モード complete で丸ごと出し直される(ここでは実行していないので件数は示さない)。 SSDSE-B-2026.csv は項目コードと日本語見出しの 2 段ヘッダなので、 そのまま置くと 2 行目の見出しがデータ行として読まれる。 置くファイルは見出し 1 行に整え、 スクリプトとして動かすなら最後に query.awaitTermination() を書かないとすぐ終了する。
Spark の 5 大事故は 「小さいデータ (数百 MB 以下) で使い性能逆転」「遅延評価で filter().show() しないと何も起きない罠」「collect() で Driver の OOM」「Python UDF が Java/Scala 比 10〜100 倍遅い」「groupBy/join の Shuffle で全ノード間ネットワーク転送」。 特に最後の Shuffle は数 TB 規模で 1 ジョブが数時間刺さる典型原因で、 broadcast join や partitionBy 設計で回避します。
df.filter(...) しても何も起きません。 .show() や .count() で初めて動く。df.collect() は全データを Driver に集める。 巨大データでメモリ爆発。groupBy, join はネットワーク経由でデータ移動が発生。 不要な shuffle を避ける設計。※ 上記は文献調査・現場経験で報告される頻度の高い注意点。 ドメインや手法のバージョンによって追加の落とし穴がある場合があります。
初学者を超えた中級者がハマる「2 周目の失敗例」を集めました。 本編の 5 件と合わせて 15 件のチェックリストとして活用してください。
.collect() は driver に全データ送付。 大規模では OOM。 take(n) や write を使う。cache() 必須。 でないと再計算で時間ロス。spark.sql.adaptive.enabled=true。spark.executor.memory=8g 以上に。本編の 10 語に加え、 さらに専門用語 20 個を整理。 文献を読むときの「分からない単語チェッカー」として使えます。
Estimator + Transformer をチェーン化するパターン。fit() でモデルを返すコンポーネント。 学習を担う。transform() でデータ変換するコンポーネント。StructType で明示。 推論より速い・確実。window 関数、 lag/lead、 rangeBetween で複雑な時系列処理可。ANALYZE TABLE で統計収集。pytest-spark で local モード単体テスト。 ETL 関数を pure function に。log4j.properties で制御。 INFO は冗長、 通常 WARN。47 都道府県 × 12 年分のデータを使い、 Apache Spark を多角的に体験する総合演習。 単発の操作ではなく、 一連の分析フローを通して理解を深めます。
🎯 このコードでやること:PySpark で SSDSE-B-2026 を全部処理
📥 入力:SSDSE-B-2026.csv の 47 都道府県データ。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 | from pyspark.sql import SparkSession from pyspark.sql.functions import sum as ssum, avg, col spark = SparkSession.builder.appName('SSDSE').getOrCreate() import pandas as pd # SSDSE-B-2026.csv は 1 行目が項目コード・2 行目が日本語見出しの 2 段ヘッダ。 # Spark の CSV 読み込みでは 1 行目を飛ばせないので、pandas で読んでから Spark に渡す pdf = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=1) df = spark.createDataFrame(pdf) df = df.withColumn('総人口', col('総人口').cast('long')) df = df.withColumn('年度', col('年度').cast('int')) # 年度別全国人口 (df.groupBy('年度') .agg(ssum('総人口').alias('全国')) .orderBy('年度') .show()) |
📤 実行結果(Spark では未実行:Spark 実行環境が無いため、 同じ集計を pandas で実測し、 show() の表の形に整えた値):
💬 結果の読み方:564 行を Spark DataFrame で集計。 同じコードが 6 億行でも動く。 これが Spark の意義。
関連概念を視覚的に整理した概念マップ。
マップ中心の Apache Spark から 6 軸が広がる。 「Hadoop (前身、 MapReduce ベース)」「Databricks (Spark 商用版)」「Delta Lake (ACID トランザクション層)」「PySpark (Python API)」「MLlib (機械学習ライブラリ)」「Structured Streaming (準リアルタイム処理)」が連結し、 TB 級データに対する分散処理 → 学習 → ストリーミング推論まで一つのフレームワークで完結する。
「Apache Spark」は クラスタ上で分散処理する大規模データ処理基盤 として、 上流の HDFS/S3 等の分散ストレージと下流の機械学習 (MLlib)・SQL 解析・ストリーミングを繋ぐ。 pandas が単一マシンで詰まる規模(数十 GB 以上)で真価を発揮する選択肢となる。
Spark は大規模データを RDD / DataFrame で分散処理する基盤で、 Hadoop / MapReduce の系譜を継ぎ、 MLlib・Structured Streaming・Delta Lake と連携してバッチ・ML・リアルタイムを一基盤で扱う。
「Spark」を使うかは、 データ量と分散環境の必要性で判断する。
SparkSQL または DuckDB/BigQuery で代替可、 No → DataFrame APISSDSE-B-2026 は 47 行 × 約 100 列で pandas で十分。 Spark はテラバイト級・分散環境が前提なので、 教育用途では概念学習に留め、 実体験は EMR/Databricks で行う。
df.filter(...) しても何も起きません。 .show() や .count() で初めて動く。df.collect() は全データを Driver に集める。 巨大データでメモリ爆発。groupBy, join はネットワーク経由でデータ移動が発生。 不要な shuffle を避ける設計。collect() で集めるwrite でファイル出力するか take(n) で先頭サンプルだけ取る。spark.sql.shuffle.partitions を core 数の 2-4 倍に。broadcast() を明示。pyspark.sql.functions のビルトインを使う。 必要なら Pandas UDF (vectorized) で 10-100倍高速化。.collect() は driver に全データ送付。 大規模では OOM。 take(n) や write を使う。cache() 必須。 でないと再計算で時間ロス。spark.sql.adaptive.enabled=true。spark.executor.memory=8g 以上に。Spark は「データエンジニアリング」分野の中で発展してきた概念・手法です。 学術的には継続的な研究で精緻化され、 実務的にはツール・ライブラリの普及で誰でも使えるようになってきました。 用語の使い方・意味は時代と分野で少しずつ変わるため、 文脈に応じた解釈が大切です。 入門書だけでなく、 標準的な教科書(例:データサイエンス・統計学の定本)や信頼できるオンライン教材も併用すると、 ぶれない理解に近づけます。
Spark は分散処理エンジンである以上、 常に問うべきは「そもそも分散が要るのか」です。 このセクションは、 本文の遊び場(遅延評価・パーティション並列)や姉妹ページ Hadoop(Map→Shuffle→Reduce)/ビッグデータとは重ならない独自の角度――規模による道具選定の線引き――に絞って掘り下げます。 題材は本コンペの実データ SSDSE-B-2026.csv です。
分散処理は「100 人で蔵書を読む」体制です。 しかし読む本が数ページなら、 100 人を招集し役割を割り振る号令の時間のほうが、 読む時間より長くなります。 分散の利得が正になるのは、 データが「1 台のメモリに載り切らない」領域に入ってからです。
SSDSE-B-2026.csv はその正反対にいます。 実測(pd.read_csv(encoding='cp932', skiprows=[1]))は次の通り:
df[df['SSDSE-B-2026']==2023] は 47 行)df.memory_usage(deep=True).sum() = 568,152 B ≒ 555 KiB、 1 行あたり約 1,007 BSpark の既定パーティションは 128 MiB。 この全データ(555 KiB)はそこに約 236 個収まり、 単独では 1 パーティション分の 0.42% しか埋めません。 つまり分割できず並列度 = 1。 分散エンジンなのに分散できず、 JVM 起動やシリアライズの純オーバーヘッドだけが乗ります。 「Spark を使う」=「速くなる」ではないことが、 実測から一目で言えます。
read_csv がミリ秒で終える 564 行を、 Spark はセッション起動だけで秒オーダー消費します。groupBy('Prefecture') ですら、 Spark は論理計画上シャッフル段を作ります。 本文の遊び場で見た「並列で割れない部分」が、 この 47 行にも生じます。 一方 pandas の .groupby('Prefecture') はインプロセスのハッシュ集約で完結し、 ネットワーク往復はゼロです。local[*] で「動いた」ことと「速い・適切」は別物です。 学習目的で Spark を触るのは有益ですが、 それを本番の道具選定の根拠にしないこと。
判断軸を定量化しましょう。 分散が勝つのは (a)単一ノードのメモリ/ディスクに収まらないか、 (b)同じジョブの反復・再利用で JVM 固定費を償却できる領域だけです。 実務的な目安:
下のスライダー(またはグラフをタップ/ドラッグ)で行数を変え、 固定費と並列利得の綱引きを見てください。 曲線は架空の説明用モデル(pandas は単一ノード線形、 Spark は「JVM/セッション固定費 ≒ 4 秒+128 MiB ごとに最大 8 コア並列」)ですが、 SSDSE-B の 564 行(赤破線)が pandas 圏に深く沈むことは実測アンカーから確かです。 交差はおよそ百万行あたりに現れます。
※ 曲線と「架空の計測点」は説明用の決定的モデル(シード付き擬似乱数)で、 実際のベンチマーク値ではありません。 一方、 564 行/112 列/351 KiB(CSV)/555 KiB(pandas メモリ)は SSDSE-B-2026.csv からの実測値です。