Data Engineer Associate — 第 4

Spark はどう動き、SQL で何をどこまで書けるか

読了 25確認 12更新 2026-08-09

この章で学ぶこと

  • 遅延評価とアクションの関係を説明する
  • Job / Stage / Task と shuffle 境界を読む
  • when と otherwise で条件分岐を書き切る
  • OVER 句の既定フレームを踏まえて累積を書く
  • 配列を高階関数で加工し Python UDF を避ける

この章に出てくる用語

遅延評価
select や filter を書いた時点では計算されず、DAG に実行計画が積まれるだけの仕組みです。show や write などのアクションを呼んで初めて Job が起動するため、存在しない列名のエラーもそこまで出ません。
シャッフル
同じキーの行を 1 か所へ集めるためのデータ再配置です。ディスク I/O とネットワーク I/O を伴い、groupBy や join で発生します。分割数の既定は spark.sql.shuffle.partitions の 200 です。
Stage
Job をシャッフルの境界で切り分けた単位です。1 つの Stage は複数の Task に分かれ、Task はパーティション 1 つにつき 1 個作られます。wide な変換が 1 回入るたびに Stage が 1 つ増えます。
when と otherwise
PySpark で行ごとに値を出し分ける式で、SQL の CASE WHEN と同じです。条件は上から評価され、最初に真になった枝が勝ちます。otherwise を省くと該当しない行は NULL になります。
ウィンドウフレーム
OVER 句の中で集計範囲を決める指定です。ORDER BY を書いてフレームを省略すると RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW が既定になり、同じ値の行がまとめて加算されます。
高階関数
配列を行にほどかず、ラムダ式で要素を加工する Spark SQL の関数群です。transform と filter は配列を返し、exists はブール値を、aggregate は 1 つの値を返します。Python UDF より高速です。
Python UDF
Python の関数を列の計算に差し込む仕組みです。1 行ごとに JVM と Python の間でデータを詰め替えるため遅く、戻り値型の既定は StringType なので、整数を返すつもりで省略すると文字列の列になります。

Spark はどう実行されるか — 遅延評価、DAG と Job / Stage / Task、narrow と wide、Driver と Executor

困りごとから入ります。ノートブックに変換を 10 行ほど書いて上から実行すると、どのセルも 1 秒たらずで終わります。ところが最後に `display(df)` と書いた 1 行だけが 40 分固まり、その間クラスターの料金は流れ続けます。原因は単純で、Spark は書いた順に処理していないからです。

遅延評価 — 書いた時点では何も走らない

Spark の操作は 2 種類しかありません。`select` `filter` `withColumn` `join` `groupBy` のような変換 (transformation) は、実行計画にメモを 1 行足すだけで、データには触れません。Apache Spark の公式ドキュメントは「Spark のすべての変換は遅延 (lazy) であり、結果をすぐには計算しない」と明記しています。計算が始まるのは `show()` `count()` `collect()` `write` `display()` といったアクション (action) を呼んだ瞬間です。

積み上がった変換のつながりを DAG (有向非巡回グラフ) と呼びます。処理どうしの依存関係を、後戻りしない矢印だけで表した図のことです。Spark UI には Job と Stage ごとに DAG Visualization という図が出ます。

遅延評価は、エラーが出るタイミングまで変えます。次のコードは存在しない列を参照していますが、例外を捕まえられません。`/Volumes` で始まるパスは、Databricks がファイルを置く場所の 1 つです。

try:
    df = spark.read.json("/Volumes/main/raw/events")
    df.withColumn("new", df["missing"] + 1)   # 変換だけ。ここでは走らない
except Exception as e:
    print("Error:", e)                        # 何も表示されない

この `try` ブロックにはアクションが 1 つも無いため、列名が実在するかの検証すら行われません。`AnalysisException` は後続の `show()` や `write` で初めて飛びます。エラーの出た行と原因の行がずれるのは Spark ではふつうのことで、「例外が出なかったから正常」と判断して本番に出すと、書き込みの瞬間に落ちます。

Job / Stage / Task — アクション 1 回が 3 段に分解される

アクションを呼ぶと、Spark は計画を 3 段に分けます。Apache Spark の用語集にある定義が、そのまま試験で問われる語彙です。

単位公式の定義個数を決めるもの
Jobアクションに応じて起動する、複数の Task からなる並列計算アクションを呼んだ回数
Stage互いに依存する小さな Task の集合へ Job を分割したものシャッフル回数 + 1
Task1 つの Executor へ送られる仕事の単位その Stage が扱うパーティション数

冒頭の「最後の 1 行だけ 40 分」は、この表から説明できます。`display(df)` がアクションなので、そこで初めて Job が 1 つ起動し、それまでの 10 行ぶんの変換がまとめて走ります。セルごとの実行時間は処理の重さを表していません。読むべきなのは Spark UI の Job と Stage の一覧です。

Task の数え方は素直で、1 パーティションにつき 1 Task です。パーティションが 8 個しかないデータを 64 コアのクラスターに投げても、同時に動くのは 8 Task で、残り 56 コアは遊びます。「クラスターを大きくしたのに速くならない」の典型がこれです。

narrow と wide — Stage の切れ目はここにできる

Stage の境界を決めるのはシャッフル (shuffle) です。公式ドキュメントは「データをパーティション間で別のグループに再配置する仕組み」と定義し、「ディスク I/O、データのシリアライズ、ネットワーク I/O を伴う高価な操作」と説明しています。

出力の 1 パーティションが入力の 1 パーティションだけを見れば作れる操作が narrow transformation、同じキーの行を 1 か所へ集めないと答えが出せない操作が wide transformation です。

操作種類シャッフル
select / filter / withColumn / castnarrow起きない
union / explode / coalescenarrow起きない
groupBy + aggwide起きる
join (ブロードキャスト以外)wide起きる
distinct / dropDuplicateswide起きる
orderBy / sort / repartitionwide起きる
ウィンドウ関数 (OVER 句)wide起きる

`coalesce(n)` は既存のパーティションを隣どうしでまとめるだけなので narrow、`repartition(n)` は全行をハッシュで振り直すので wide です。出力ファイル数を減らしたいだけなら `coalesce` のほうが安く済みます

書く順番にも効きます。行と列を削る操作どうしの前後は Catalyst (記述を実行計画へ組み立て直す最適化エンジン) が自動で入れ替えますが、集約をまたぐ絞り込みは前倒しできません。`groupBy().agg().filter()` は `HAVING` 相当で、集約前の絞り込みとは結果そのものが違うからです。

この値を知らないと何が起きるか

設定既定値効き方
spark.sql.shuffle.partitions200シャッフル後のパーティション数、つまり次の Stage の Task 数
spark.sql.files.maxPartitionBytes134217728 (128 MB)ファイル読み込み時の 1 パーティションの上限
spark.sql.autoBroadcastJoinThreshold10485760 (10 MB)これ以下なら小さい側を全 Executor へ配り join のシャッフルを回避
spark.databricks.optimizer.adaptive.enabledtrueDatabricks で AQE を有効にする

200 という固定値は両側に外れます。数 MB のテーブルを `groupBy` しても 200 Task が作られ、1 Task が数十 KB を処理して終わる空回りになります。逆に 1 TB を 200 分割すると 1 パーティション 5 GB になり、メモリに載らず spill (ディスクへの書き出し) が起きます。Databricks は固定値ではなく `spark.sql.shuffle.partitions=auto` を推奨しています。読み取り側の 128 MB からは、20 GB の Parquet (列ごとに固めて保存する分析向けのファイル形式) が約 160 パーティションになると見積もれます。

実行の途中で計画を作り直すのが AQE (Adaptive Query Execution) です。シャッフルとブロードキャストの交換が終わった時点 (AQE はこの区切りを query stage と呼びます) で実測の統計を読み、細かすぎるパーティションの結合、両側を並べ替えて突き合わせる sort merge join から小さい側を配る broadcast hash join への切り替え、偏ったパーティションの分割、空の関係の検出という 4 つの手を打ちます。論理プランの最適化そのものは Catalyst の担当で、この章では扱いません。

Driver と Executor — 結果がどちらのメモリに載るか

Driver は main() を実行して SparkContext を作るプロセスExecutor は worker node 上で Task を実行しデータを保持するプロセスです。`collect()` や `toPandas()` が危険なのはここに理由があり、全行を Driver 1 台のメモリへ集めるため、Executor をいくら増やしても Driver だけが落ちます。計画を手元で確かめるには `df.explain()` を使います。mode は simple / extended / codegen / cost / formatted の 5 種類で、simple は物理プランのみ、extended は論理プランと物理プランの両方を表示します。

確認 — 穴あけ 3

0 / 3

空欄を押すと選択肢が出ます。間違えても減点はありません。

シャッフル後のパーティション数を決める spark.sql.shuffle.partitions の既定値は です。

Job をシャッフルの境界で切り分けた単位を と呼びます。

変換だけを書いたコードでは列名の誤りが検出されず、 を呼んだ時点で AnalysisException が出ます。

Spark SQL の中核 — SELECT / JOIN / CTE / EXPLODE / PIVOT とビューの使い分け

困りごとから入ります。要件を渡されたとき毎回迷うのが「これは SQL 1 本で書けるのか、Python に落として組み立てるのか」という線引きです。迷った末に Python の `for` ループと自作関数で書き、レビューで「それ SQL で書けます」と言われる。この節は構文の一覧表ではなく、SQL でどこまで書けるかの境界線を引くことに使います。

まず前提 — SQL と DataFrame API に性能差はない

Databricks で `spark.sql("SELECT ...")` と書いても、`df.groupBy(...).agg(...)` と書いても、最終的には同じ論理プランに変換されます。どちらが速いという話ではなく、読みやすいほうを選んでよいのが Spark SQL の設計です。境界が生まれるのは、Python の関数を UDF (ユーザー定義関数) として差し込んだ瞬間です。UDF の中身は Spark から見えない箱になり、最適化の対象から外れます。つまり「SQL で書ける範囲」は「最適化してもらえる範囲」と同義です。

行ごとの条件分岐は CASE WHEN で書き切れる

UDF に逃げがちな処理の筆頭が条件分岐ですが、これは SQL 側にあります。SQL では `CASE WHEN 条件 THEN 値 ... ELSE 既定値 END`、PySpark では `when(条件, 値).otherwise(既定値)` で、両者は同じものです。`.when()` を並べれば if / elif / else と同じ構造になります。

df.withColumn("category",
    when(col("price") > 1000, "High")
    .when(col("price") > 300, "Mid")
    .otherwise("Low"))

押さえるべき仕様が 4 つあります。1 つ目、条件は上から順に評価され、最初に真になった枝の値だけが採用されます。上のコードで price=1200 は 1 つ目で `High`、price=500 は 1 つ目が偽で 2 つ目が真なので `Mid`、price=100 はどちらも偽で `Low` です。2 つの `when` を逆順に書くと、1200 も `Mid` で止まります。

2 つ目、`ELSE` `otherwise` は省略できますが、省略すると既定は NULL です。Databricks SQL のリファレンスは「`def` を省略した場合の既定は NULL」明記しています。どの条件にも当たらない行が NULL になり、後段の `filter` で静かに消えます。

3 つ目、比較演算子の境界です。`col("amount") > 1000` は等号を含まないので、amount がちょうど 1000 の行は偽の側に落ちます。`>=` との取り違えは構文エラーにならず、件数だけがずれます。4 つ目、`withColumn(名前, 式)` は同名の列があれば上書き、なければ追加です。`df.withColumn("price", col("price") * 1.1)` は `price_1` のような別列を作らず、`price` を置き換えます。行の絞り込みは `filter`、列の追加は `withColumn`、集約は `groupBy().agg()` と役割が分かれている点も、そのまま出題されます。

SELECT と GROUP BY — 並べてよい列は 2 種類だけ

`GROUP BY` を書いたクエリの `SELECT` と `ORDER BY` に並べてよいのは、(1) `GROUP BY` に含めた列と、(2) 集計関数の結果、の 2 種類だけです。それ以外を書くと `MISSING_AGGREGATION` 系のエラーになります。次の SQL が落ちるのは `age` がどちらにも該当しないからです。

SELECT department, COUNT(*) AS cnt
FROM employees
GROUP BY department
ORDER BY age;          -- age は GROUP BY にも集計にも無い

直し方は `GROUP BY department, age` にするか、`ORDER BY cnt DESC` に変えるかの二択です。`SELECT` で付けた別名は `ORDER BY` から参照できます。集計後に絞り込むなら `HAVING`、集計前なら `WHERE` で、DataFrame API では `groupBy().agg().filter()` の最後の `filter` が `HAVING` に当たります。

JOIN — 行が減る原因はほぼ NULL と型

INNER JOIN の結果が想定より大幅に少ないとき、まず疑うのは結合キーです。SQL の三値論理では `NULL = NULL` が真にならないため、キーが NULL の行は静かに全部落ちます。エラーも警告も出ないので、件数を数えるまで気づけません。前段で `WHERE key IS NOT NULL` を入れて落ちる行を可視化し、必要なら別テーブルへ退避します。`coalesce(key, 0)` のような穴埋めは、本来無関係な行どうしを結合させる事故になります。

もう 1 つの原因が型と表記のずれで、`STRING` と `INT` の暗黙キャスト、前後の空白、大文字小文字の違いは、いずれも一致しない側に倒れます。なお「右に存在するかどうかだけ知りたい」なら `LEFT SEMI JOIN`、「存在しないものを取りたい」なら `LEFT ANTI JOIN` を使うと、右側の列を返さず行数も増えません。10 億行と 1 万行のような組み合わせでは第 1 節の Broadcast Join が効き、SQL では `/*+ BROADCAST(regions) */`、PySpark では `broadcast(regions)` で明示できます。

EXPLODE と PIVOT — 形を変える 2 つの操作

JSON から取り込んだテーブルには、1 行の中に配列が入っている列がよくあります。`explode` はこれを行へほどく関数で、ARRAY を渡すと要素ごとに 1 行を作り列名は `col`、MAP を渡すと `key` と `value` の 2 列になります。ここで入力が NULL または空配列だと 1 行も出力されず、その行ごと結果から消えます。行を残したいときは `explode_outer`、要素の位置も欲しいときは `posexplode` です。Databricks Runtime 12.2 以降では `SELECT` リストや `LATERAL VIEW` からの呼び出しは非推奨で、`FROM` 句のテーブル参照として書くのが推奨です。

逆方向の `PIVOT` は、ある列の値を列名へ展開してクロス集計表を作ります。構文は `PIVOT (集計 FOR 列 IN (値, ...))` で、`IN` に並べる値はリテラルでなければならず、`1 + 0` のような式を書くと `NON_LITERAL_PIVOT_VALUES` になります。集計を複数書いた場合、出力列名は `q1_total` のように「値の別名_集計の別名」になります。列を行へ戻す操作は `UNPIVOT` で、配列を行へ開く `explode` とは目的が違います。

名前を付ける道具と、意外な既定値

`WITH 名前 AS (クエリ)` の CTE は読みやすさのための道具で、結果を保存する仕組みではありません。同じ CTE を 2 か所から参照すれば 2 回評価されることがあります。再帰 CTE は Databricks Runtime 17.0 以降で使え、`MAX RECURSION LEVEL` の既定は 100、結果セットは `LIMIT` が無い限り 1,000,000 行に制限されます。ビューは 3 種類で、`CREATE VIEW` はカタログに永続、`CREATE TEMPORARY VIEW` は作成セッション限り、`CREATE GLOBAL TEMPORARY VIEW` は `global_temp` スキーマ経由で他セッションからも見えます。いずれもクエリ定義の保存であって、データのコピーではありません。

集合演算では既定値が事故になります。Databricks SQL の `UNION` `INTERSECT` `EXCEPT` は既定が `DISTINCT` で、重複行が黙って消えます。単純に連結したいだけなら `UNION ALL` を明示します。最後に `QUALIFY` です。`WHERE` が集計前、`HAVING` が集計後を絞るのに対し、`QUALIFY` はウィンドウ関数の計算後を絞ります。次の節で扱う「グループごとに 1 件だけ取る」が、サブクエリ無しの 1 文で書けます。

SELECT customer_id, order_id, order_date
FROM orders
QUALIFY ROW_NUMBER() OVER (
  PARTITION BY customer_id ORDER BY order_date DESC) = 1;

確認 — 穴あけ 3

0 / 3

空欄を押すと選択肢が出ます。間違えても減点はありません。

when チェインで otherwise を省くと、どの条件にも当たらなかった行の値は になります。

withColumn に既存と同じ列名を渡すと、その列は されます。

GROUP BY のあるクエリの SELECT と ORDER BY に並べてよいのは、GROUP BY のキーと の結果だけです。

ウィンドウ関数とフレーム指定 — OVER 句、rowsBetween と rangeBetween、row_number と rank の違い

困りごとから入ります。「顧客ごとに最新の注文を 1 件だけ取ってきて」という依頼は、`GROUP BY customer_id` では解けません。`MAX(order_date)` は取れますが、一緒に欲しい商品名や金額を `SELECT` に並べた瞬間に第 2 節のエラーになります。行を減らさずに集計値や順位を付けたい——これがウィンドウ関数の出番です。

OVER 句 — 行数を保ったまま計算する

`GROUP BY` はグループごとに 1 行へ潰します。ウィンドウ関数は元の行数を保ったまま、各行に集計値や順位を列として足します。書き方は集計関数のうしろに `OVER` 句を付けるだけです。

SELECT customer_id, order_date, amount,
       SUM(amount) OVER (PARTITION BY customer_id ORDER BY order_date) AS running_total
FROM orders;

`OVER` の中身は 3 つのパーツでできています。`PARTITION BY` が計算を区切る単位、`ORDER BY` がその中での並び、そしてフレーム句が「同じ区切りの中でもどこからどこまでを足すか」です。3 つとも省略でき、省略したときの挙動を知らないことが事故の元になります。PySpark では同じものを `Window` クラスで組み立てます。

from pyspark.sql.window import Window
from pyspark.sql import functions as F

w = Window.partitionBy("customer_id").orderBy("order_date") \
          .rowsBetween(Window.unboundedPreceding, Window.currentRow)
df2 = df.withColumn("running_total", F.sum("amount").over(w))

境界の定数は `Window.unboundedPreceding` (区画の先頭)、`Window.currentRow` (現在行)、`Window.unboundedFollowing` (区画の末尾) の 3 つです。行ベースのオフセットは現在行を 0 とした相対位置なので、`rowsBetween(Window.unboundedPreceding, 0)` は上のコードとまったく同じ累積フレームを指します。試験では整数リテラルで書かれた選択肢も出るので、`0` と `Window.currentRow` は同じものだと結び付けておきます。同様に `-2` は 2 行前です。

既定のフレーム — 累積和がずれる最大の原因

フレームを書かなかったときの既定は、`ORDER BY` の有無で変わります。Databricks の `Window` クラスのドキュメントは、順序が定義されていない場合は無制限フレーム (rowFrame, unboundedPreceding, unboundedFollowing)順序が定義されている場合は伸びていくフレーム (rangeFrame, unboundedPreceding, currentRow) が既定だと書いています。SQL に直すと、前者が `ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING`、後者が `RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW` です。

問題になるのは `RANGE` のほうです。`ORDER BY order_date` だけを書いて累積和を取ると、既定は `RANGE` なので同じ日付の行がすべて 1 つの塊として同時に加算されます。1 日に 3 件の注文がある顧客では、3 行とも同じ累積値が並び、1 件ずつ増える階段になりません。意図が「行単位の累積」なら `ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW` を明示するのが正解です。

rowsBetween と rangeBetween — 数えるものが違う

2 つの違い一言で言えます。`ROWS` は行数を数え、`RANGE` `ORDER BY` 式の値を数えます。給与が同額の 2 人を並べたとき、ROWS 1 行ずつ数えるので累計が違い、RANGE同じ値を共有するので累計が同じになります。

やりたいことフレーム指定 (PySpark)
累積 (running total)rowsBetween(Window.unboundedPreceding, Window.currentRow)
逆累積 (この行から末尾まで)rowsBetween(Window.currentRow, Window.unboundedFollowing)
直近 3 行の移動合計rowsBetween(-2, Window.currentRow)
区画全体の合計を各行に付けるrowsBetween(Window.unboundedPreceding, Window.unboundedFollowing)

`RANGE` には制約もあります。`ORDER BY` が無いと `DATATYPE_MISMATCH.RANGE_FRAME_WITHOUT_ORDER`、`ORDER BY` に式が 2 つ以上あると `DATATYPE_MISMATCH.RANGE_FRAME_MULTI_ORDER` のエラーになります。オフセットは `ORDER BY` 式から足し引きできる型でなければならず、日付列なら `INTERVAL` を渡します。

row_number と rank と dense_rank — 同点の扱いだけが違う

冒頭の「顧客ごとに最新 1 件」に戻ります。答えは `row_number()` です。

SELECT * FROM (
  SELECT *, ROW_NUMBER() OVER (PARTITION BY customer_id ORDER BY order_date DESC) AS rn
  FROM orders
) WHERE rn = 1;
関数同点の扱い並び (1, 1, 2 の入力)= 1 で 1 行に絞れるか
row_number必ず別々の番号1, 2, 3絞れる
rank同じ番号を共有し次を飛ばす1, 1, 3絞れない
dense_rank同じ番号を共有し飛ばさない1, 1, 2絞れない

ここで `RANK()` を選ぶと、同じ日付の注文が 2 件ある顧客で `= 1` が 2 行返り、「1 件に絞る」という要件を満たせません。順位関数は `ORDER BY` が必須で、フレーム句は書けません。第 2 節の `QUALIFY` を併用すれば、この 2 段のサブクエリを 1 文に畳めます。

省略したときに壊れるもの

`ORDER BY` を省くと、`row_number()` や `lag()` の結果は非決定的になります。同じ入力でも実行のたびに値が変わるため、日次バッチで前日と差分を取ると毎回差分が出ます。並び順が一意に決まる列 (タイムスタンプ + ID など) を必ず指定します。

`PARTITION BY` を省くと、全行が 1 つのウィンドウにまとめられます。結果が「カテゴリ別ランキング」ではなく「全件ランキング」になり、「ユーザーごとの累計のはずが全行に総合計が入っている」という症状になります。実行面ではもっと深刻で、全行が 1 つのパーティションへ集められ 1 Task で処理されます。ウィンドウ関数は第 1 節でいう wide transformation で、`PARTITION BY` のキーでシャッフルが起きます。キーの値が偏れば特定の Task だけが巨大になり、spill や OOM (メモリ不足) の原因になります。AQE の偏り対策は結合が対象なので、ウィンドウ側の偏りはキー設計で避けるしかありません。

前後の行を参照する `LAG` と `LEAD` は、`ORDER BY` の順で 1 行ずらした値を返し、端の行では NULL になります。これらと `FIRST_VALUE` `LAST_VALUE` には `IGNORE NULLS` を付けられ、NULL を飛ばして「直前の有効な計測値」を引けます。ここで `collect_list` で全明細を 1 配列に集めて Python で計算したり、`crossJoin` で全組み合わせを作って累計を出したりするのは、どちらも行数とメモリが破綻する典型的なアンチパターンです。

確認 — 穴あけ 3

0 / 3

空欄を押すと選択肢が出ます。間違えても減点はありません。

ORDER BY を書いてフレームを省略したときの既定は BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW です。

顧客ごとに最新の 1 件だけを取り出すには を使い、= 1 で絞ります。

ウィンドウ関数で を省くと、全行が 1 つのウィンドウにまとめられ単一 Task に集中します。

複合型と高階関数 — collect_list / collect_set と transform / filter / exists / aggregate

困りごとから入ります。1 行の中に配列が入っているテーブルを渡されたとき、多くの人がやるのは「`explode` で行にほどいて、加工して、`collect_list` で配列に戻す」という往復です。書けはしますが、行を増やして減らす操作なのでシャッフルが 2 回入ります。Python の `for` ループを UDF にして配列を回す手もありますが、こちらは 1 行ごとに JVM と Python の間でデータを詰め替えるので、さらに遅くなります。配列は配列のまま、行の中で加工する——それが高階関数です。

collect_list と collect_set — 行を配列にたたむ

まず配列を作る側です。`collect_list(col)` は `GROUP BY` と組み合わせて、グループ内の値を `ARRAY` にまとめる集約関数で、Databricks では `array_agg` の別名です。重複を落としたいときは `collect_set(col)`、あるいは `collect_list(DISTINCT col)` と書きます (`DISTINCT` を付けると `collect_set` と同義になります)。

SELECT customer_id, collect_list(product_id) AS products
FROM orders
GROUP BY customer_id;

仕様で押さえるべき点が 2 つあります。1 つ目、NULL は要素に含まれません。公式ドキュメントは「NULL 値は除外される」と明記しています。したがってグループの値がすべて NULL のときの結果は NULL ではなく空の配列です。「欠損も含めて数えたい」なら `collect_list(coalesce(col, 'unknown'))` のように埋めてから集約するか、事前に `WHERE` で不正値を除外します。

2 つ目、配列の中の順序は非決定的です。パーティションの並びに依存するため、上流で `orderBy` してもシャッフルをまたいだ時点で保証されません。順序が意味を持つなら `array_sort` や `sort_array` を必ずかぶせます。

失敗としてよくあるのが 2 つあります。1 つは集約キーの取り違えで、「ユーザーごとの購入商品リスト」を作るつもりで `groupBy("shop_id")` と書いても構文エラーは出ず、店舗ごとに全ユーザーの商品が混ざった配列が静かに出来上がります。そもそも `groupBy` を付けずに呼べば、データセット全体が 1 つの配列になります。もう 1 つは、すでに配列型の列にもう一度 `collect_list` をかけて `ARRAY<ARRAY<STRING>>` という二重配列ができるパターンです。中間テーブルで集約済みの列を別の `groupBy` でもう一度たたむと起きます。`flatten` で平らにするか、集約を 1 回に整理します。

transform と filter — 配列の形を保ったまま加工する

`transform(expr, func)` は配列の各要素にラムダ式を適用し、同じ長さの新しい配列を返します。ラムダには 2 つの形があり、要素だけを受ける `x -> x + 1` と、要素と 0 始まりの添字を受ける `(x, i) -> x + i` が書けます。`transform(array(1, 2, 3), x -> x + 1)` は `[2,3,4]`、`transform(array(1, 2, 3), (x, i) -> x + i)` は `[1,3,5]` になります。

`filter(expr, func)` は述語が真の要素だけを残し、`expr` と同じ型の配列を返します。`filter(array(1, 2, 3), x -> x % 2 == 1)` は `[1,3]`、`filter(array(0, null, 2, 3, null), x -> x IS NOT NULL)` は `[0,2,3]` です。NULL を落とす処理も、専用関数ではなくこのラムダで書きます。

ここを取り違える出題が多いので整理します。長さ 3 以上の要素だけを残したいのに `transform(tags, t -> length(t) >= 3)` と書くと、返るのは元の文字列ではなくブール値の配列です。要素を選ぶなら `filter`、要素を作り替えるなら `transform`、という対応で覚えます。

exists と aggregate — 配列を 1 つの値へ落とす

`exists(expr, func)` は「条件を満たす要素が 1 つでもあるか」を返します。落とし穴は三値論理です。条件を満たす要素があれば `true`、無ければ `false` ですが、満たす要素が無く、かつ配列に NULL が含まれる場合は `NULL` を返します。公式ドキュメントの例では `exists(array(1, NULL, 3), x -> x % 2 == 0)` が `NULL` です。これを `WHERE` に置くと、その行は真ではないので結果から消えます。「条件に合わない行を数えたら合計が合わない」の正体がこれです。全要素が条件を満たすかを見たいときは `forall` を使います。

`aggregate(expr, start, merge [, finish])` は配列を 1 つの値へ畳み込む関数で、`reduce` の別名です。`start` が初期値、`merge` が要素を足し込むラムダ、省略可能な `finish` が最後の仕上げです。

SELECT aggregate(array(1, 2, 3), 0, (acc, x) -> acc + x);                  -- 6
SELECT aggregate(array(1, 2, 3), 0, (acc, x) -> acc + x, acc -> acc * 10);  -- 60

注意すべきは型です。`finish` を省いた場合、結果の型は `start` の型に一致します。平均を求めるつもりで `start` に整数の `0` を渡すと、途中の計算まで整数に丸められます。小数が必要なら初期値を `0D` のように書くか、`finish` で割り算します。

なぜ UDF ではなく高階関数なのか

ここまでの関数はすべて Spark SQL の式として表現されるため、Catalyst の最適化とコード生成の対象になり、Photon (Databricks の C++ 実行エンジン) でも動きます。一方 Python UDF は、Executor の JVM から Python プロセスへ 1 行ずつデータを送り、結果を戻す往復が入ります。Spark からは中身の見えない箱なので、述語の押し込みも列の刈り込みも効きません。どうしても Python のロジックが必要なら、Apache Arrow でまとめて転送する Pandas UDF が次善の策で、組み込み関数と高階関数 → Pandas UDF → Python UDF という優先順位がそのまま設計の指針になります。

UDF を書くときは型の指定でつまずきます。`pyspark.sql.functions.udf(f=None, returnType=StringType(), useArrow=None)` という署名のとおり、`returnType` の既定は `StringType` です。整数を返すつもりで型を省くと、例外は出ないまま文字列の列ができあがり、下流の計算で崩れます。`udf(lambda x: x + 1, IntegerType())` のように必ず明示します。Pandas UDF も同じで、`@pandas_udf("long")` のように戻り値の型を宣言する必要があり、Spark が自動で推論してくれることはありません。

確認 — 穴あけ 3

0 / 3

空欄を押すと選択肢が出ます。間違えても減点はありません。

配列の各要素を別の値へ作り替え、同じ長さの配列を返す高階関数は です。

collect_list は NULL を要素に含めないため、グループの値がすべて NULL なら結果は になります。

pyspark.sql.functions.udf で returnType を省略すると、既定の が使われます。

この章のまとめ

  1. アクションを呼ぶまで変換は走らず、エラーも表面化しない
  2. wide な変換が shuffle を起こし、そこで Stage が切れる
  3. otherwise を省くと NULL、withColumn は同名列を上書きする
  4. ORDER BY だけのウィンドウは RANGE 既定で同値行がまとまる
  5. 配列は explode で戻さず transform や filter で加工する

この端末にだけ保存されます(登録不要)

この章の根拠

最終確認 2026-08-09 / 対応バージョン DEA 2026-05-04 改訂版

コース全体

  1. Databricks とレイクハウスの全体像 — なぜ必要か、何がどこに属し、どの計算資源で動くか
  2. データはどこに、どんな形で保存されるか — テーブルの実体と Delta Lake の内部
  3. 計算資源を作り込み、コードを書く場所を決める — クラスタ設定とノートブック/ローカル IDE
  4. Spark はどう動き、SQL で何をどこまで書けるか
  5. PySpark でデータを加工し、遅いコードを見抜く
  6. 流れ込むデータを受け止める — Structured Streaming と Auto Loader のファイル検出
  7. 取り込み時にスキーマをどう扱うか — Auto Loader のスキーマ進化と COPY INTO
  8. ファイル以外からも取り込む — read_files・Lakeflow Connect・JDBC と API
  9. 生データを使える形に育てる — メダリオン設計と宣言的パイプライン
  10. 処理を1つのジョブに束ねる — タスク種別と依存関係の設計
  11. ジョブを動かし、失敗から立て直す — トリガー・リトライ・パラメータ
  12. 書いたものを安全に本番へ届ける — Git 連携とデプロイの自動化
  13. 誰に何を見せるか — Unity Catalog のアクセス制御とアイデンティティ統制
  14. 誰が何をしたか、いくらかかったか、どこで詰まったかを見る
  15. テーブルを保守して速くする — 運用コマンド・保持期間・レイアウト設計・障害復旧
  16. DEA 直前仕上げ — 方式選定・出題範囲対応表・数値総まとめ・引っかけの型