テーマ切替
Professional 学習教材 ① Python/SQL でのデータ処理コード開発(配点 22%)
Databricks 公式ドキュメント(日本語版)の内容を、重要用語・概念を漏らさずまとめた Professional 試験向け自習教材です。 参照した公式ページ:
- Azure Databricks における PySpark: https://learn.microsoft.com/ja-jp/azure/databricks/pyspark/
- チュートリアル: Apache Spark DataFrame を使用してデータを読み込んで変換する: https://learn.microsoft.com/ja-jp/azure/databricks/getting-started/dataframes
- ユーザー定義関数 (UDF) とは: https://learn.microsoft.com/ja-jp/azure/databricks/udf/
- Pandas のユーザー定義関数: https://learn.microsoft.com/ja-jp/azure/databricks/udf/pandas
- ユーザー定義スカラー関数 - Python: https://learn.microsoft.com/ja-jp/azure/databricks/udf/python
- ウィンドウ関数: https://learn.microsoft.com/ja-jp/azure/databricks/sql/language-manual/sql-ref-window-functions
- ウィンドウ フレーム句: https://learn.microsoft.com/ja-jp/azure/databricks/sql/language-manual/sql-ref-syntax-window-functions-frame
- 開発言語の選択: https://learn.microsoft.com/ja-jp/azure/databricks/languages/overview
- 上位の関数(高階関数): https://learn.microsoft.com/ja-jp/azure/databricks/semi-structured/higher-order-functions
注: 高階関数ページ(semi-structured/higher-order-functions)は本文が短く、詳細は参照ノートブックへのリンク中心でした。3-4 節の関数個別解説は Apache Spark SQL の標準関数仕様を補足しています。
1. このドメインの概要
このドメイン「Python/SQL でのデータ処理コード開発」は Professional 試験で配点 22%・最重要の領域です。単なる用語暗記ではなく、PySpark / Spark SQL のコードを読んで挙動・結果・性能を答えられることが求められます。
Professional で問われる中心テーマは次の通りです。
- DataFrame API の内部モデル(変換とアクション、遅延評価、不変性、実行計画の最適化)
- 主要な DataFrame 操作(
select/filter/withColumn/join/groupBy/agg/unionなど)を組み合わせたコード読解 - ウィンドウ関数(
PARTITION BY/ORDER BY/ フレーム指定ROWS・RANGE)による累積・移動平均・ランキング計算 - 複雑・ネストしたデータ(
struct/array/map、explode、高階関数)の処理 - UDF の種類と性能特性(組み込み関数 > SQL UDF > Scala UDF > pandas UDF > Python UDF の性能順、Unity Catalog 管理 vs セッションスコープ)
- 言語の使い分け(Python と SQL 推奨、Scala / R は制限あり)
Databricks の基盤:
- Azure Databricks は、ビッグデータと機械学習向けの統合分析エンジンである Apache Spark 上に構築されている。
- PySpark は、Python を使って Apache Spark とインターフェイスするための API。Python と Apache Spark の機能を組み合わせたもの。
- DataFrame は Apache Spark の主となるオブジェクトで、Resilient Distributed Datasets (RDD) を基盤に構築された抽象化。
- Spark DataFrames と Spark SQL は統合された計画・最適化エンジン(Catalyst)を使うため、Python / SQL / Scala / R のいずれでもほぼ同じパフォーマンスが得られる。
2. 重要用語集
| 用語(日本語) | English | 説明 |
|---|---|---|
| Apache Spark | Apache Spark | ビッグデータ・機械学習向けの統合分析エンジン。Azure Databricks の基盤。 |
| PySpark | PySpark | Python から Apache Spark を操作する API。Python と Spark の機能を組み合わせたもの。 |
| データフレーム | DataFrame | 名前付き列に編成された 2 次元のラベル付きデータ構造。スプレッドシート/SQL テーブルに相当。Spark の主オブジェクト。 |
| RDD(回復力のある分散データセット) | Resilient Distributed Dataset (RDD) | DataFrame の基盤となる低レベルの分散データ抽象化。DataFrame は RDD 上に構築される。 |
| スキーマ | Schema | DataFrame の列の名前と型の定義。手動定義、データソースからの読み取り、推論(inferSchema)が可能。 |
| 行 | Row | DataFrame 内のレコード。Row オブジェクトとして表される。最適化のためキャッシュとシャッフルに使われる。 |
| 列 | Column | スプレッドシートの列に相当。単純型(文字列・整数)だけでなく配列・マップ・null などの複雑型も表す。 |
| 変換 | Transformation | データを読み込み・操作する手順(spark.read / join / 集計 / 型キャストなど)。遅延評価され、新しい DataFrame を返す。 |
| アクション | Action | 変換の結果の計算をトリガーし値を返す操作(display / show / take / count / saveAsTable など)。 |
| 遅延評価(遅延読み込み) | Lazy Evaluation / Lazy Loading | 変換はアクションが呼ばれるまで実行されない。論理クエリを一連の命令として保持し、最も効率的な物理プランを特定してから実行。 |
| 即時実行 | Eager Execution | pandas DataFrame のモデル。定義した時点で即座に実行される。Spark の遅延評価とは対照的。 |
| 不変性 | Immutability | DataFrame は不変。変換は元を変更せず新しい DataFrame を返すため、結果は変数に保存する必要がある。 |
| 実行計画 / 物理プラン | Execution Plan / Physical Plan | Spark が変換ロジックを評価するために特定する最も効率的な実行手順。Catalyst オプティマイザが生成。 |
| 狭い変換 | Narrow Transformation | 各出力パーティションが単一の入力パーティションのみに依存(select / filter / withColumn など)。シャッフル不要で高速。 |
| 広い変換 | Wide Transformation | 出力パーティションが複数の入力パーティションに依存(groupBy / join / orderBy / distinct など)。シャッフルを伴う。 |
| シャッフル | Shuffle | パーティション間でのデータの再分配。ネットワーク I/O が発生する高コスト操作。ステージ境界を作る。 |
| Spark SQL | Spark SQL | SQL クエリと Spark プログラムを混在できるリレーショナル処理モジュール。 |
| SparkSession | SparkSession | Spark 機能へのエントリポイント。ノートブックでは spark として利用可能。 |
| ウィンドウ関数 | Window Function | ウィンドウ(関連する行のグループ)に基づいて各行の値を計算する関数。集計と異なり行数を減らさない。 |
| OVER 句 | OVER clause | ウィンドウ関数にウィンドウ仕様(window_spec)を関連付ける句。 |
| パーティション分割 | PARTITION BY | ウィンドウ関数が作用する行のグループを定義。省略時は全行が 1 パーティション。 |
| 並べ替え | ORDER BY | パーティション内の行の順序を指定。ランキング関数・分析関数・RANGE フレームに必須。 |
| ウィンドウフレーム | Window Frame | パーティション内で集計/分析関数が作用する行のスライディングサブセット。ROWS または RANGE で指定。 |
| ROWS フレーム | ROWS frame | 現在の行の前後の物理的な行数でフレーム境界を指定。 |
| RANGE フレーム | RANGE frame | 現在の行の ORDER BY 値からの値のオフセットでフレーム境界を指定。ORDER BY 式は 1 つのみ必須。 |
| ランキング関数 | Ranking Function | ROW_NUMBER / RANK / DENSE_RANK / NTILE / PERCENT_RANK。ORDER BY 必須、フレーム指定不可。 |
| 分析関数 | Analytic Function | LAG / LEAD / FIRST / LAST / NTH_VALUE / CUME_DIST。前後の行の値を参照。 |
| 集計関数 | Aggregate Function | SUM / AVG / COUNT / MIN / MAX。ウィンドウと組み合わせて累積・移動計算に使う。 |
| 構造体 | struct | 名前付きフィールドを持つネストしたレコード型。col.field でアクセス。 |
| 配列 | array | 同じ型の要素の順序付きコレクション。explode で行に展開可能。 |
| マップ | map | キーと値のペアのコレクション。 |
| explode | explode | 配列/マップの各要素を個別の行に展開する関数。 |
| 高階関数 | Higher-Order Function | 配列を受け取り、ラムダ関数を各要素に適用する関数(transform / filter / aggregate / exists など)。 |
| ラムダ関数(匿名関数) | Lambda / Anonymous Function | 高階関数に渡す、要素ごとの処理を定義する無名関数。 |
| ユーザー定義関数 | User-Defined Function (UDF) | 組み込み関数で表現困難なカスタムロジックを再利用・共有するための関数。 |
| スカラー UDF | Scalar UDF | 1 行を処理し、各行に 1 つの結果値を返す UDF。 |
| バッチスカラー UDF | Batch Scalar UDF | 1:1 の入出力を維持しつつ行をバッチ処理する UDF。行ごとのオーバーヘッドを軽減。 |
| pandas UDF(ベクトル化 UDF) | pandas UDF / Vectorized UDF | Apache Arrow でデータ転送し pandas で処理。Python UDF より最大 100 倍高速。 |
| UDTF(ユーザー定義テーブル関数) | User-Defined Table Function (UDTF) | 入力行ごとに複数の行(複数の列)を返す UDF。 |
| UDAF(ユーザー定義集計関数) | User-Defined Aggregate Function (UDAF) | 複数行を処理し 1 つの集計結果を返す UDF。セッションスコープ限定。 |
| Apache Arrow | Apache Arrow | 言語をまたぐ列指向インメモリ形式。pandas UDF でシリアライズコストを削減。 |
| Unity Catalog 管理 UDF | Unity Catalog-governed UDF | UC に永続化され、権限管理・共有・検出が可能な UDF。 |
| セッションスコープ UDF | Session-scoped UDF | 現在の SparkSession に限定される一時的な UDF。 |
| シリアライズのオーバーヘッド | Serialization Overhead | データを JVM ⇔ Python インタプリタ間で移動する際のコスト。Python UDF が遅い主因。 |
| 一時ビュー | Temporary View / Temp View | createOrReplaceTempView で作成する、言語間でデータを共有できるセッションスコープのビュー。 |
| マジックコマンド | Magic Command | ノートブックのセル言語を切り替える %python / %sql / %scala / %r / %md。 |
3. 詳細解説
3-1. PySpark と DataFrame API の基礎(変換・アクション・遅延評価・実行計画)
DataFrame とは
- DataFrame = 名前付き列に編成された、潜在的に異なる型の列を持つ 2 次元のラベル付きデータ構造。スプレッドシート、SQL テーブル、オブジェクトの辞書のようなもの。
- 列選択・フィルター・結合・集計などの豊富な機能セットを持つ。
- RDD(Resilient Distributed Datasets)を基盤に構築された抽象化。
- Spark DataFrames / Spark SQL は統合された計画・最適化エンジンを使うため、Python / SQL / Scala / R でほぼ同じパフォーマンスになる。
DataFrame の重要な要素:
- スキーマ(Schema): 列の名前と型を定義。手動定義、データソースからの読み取り、推論(
inferSchema)が可能。データ形式によってスキーマの適用セマンティクスは異なる。 - 行(Row): レコードは
Rowオブジェクトとして表される。基盤の Delta Lake などは列でデータを格納するが、Spark は最適化のため行を使ってキャッシュとシャッフルを行う。 - 列(Column): 単純型だけでなく配列・マップ・null などの複雑型も表せる。
.dropやselectでの省略で結果から除かれるだけで、データソースから列が物理的に削除されることはない。
変換(Transformation)とアクション(Action)— Professional 最重要概念
Spark は**遅延評価(lazy evaluation)**で変換とアクションを処理する。
変換(Transformation)
- 処理ロジックを表現する手順。DataFrame を返す。
- 一般的な変換: データ読み取り(
spark.read/spark.table)、結合、集計、型キャスト。 - すぐには実行されない。論理クエリを「データソースに対する一連の命令」として格納する(メモリ内の結果としてではない)。
アクション(Action)
- 1 つ以上の DataFrame で一連の変換の結果を計算するよう Spark に指示する。値を返す。
- 種類:
- 出力するアクション:
display、show - データを収集するアクション(
Rowを返す):take(n)、first、head - データソースに書き込むアクション:
saveAsTable - 計算をトリガーする集計:
count
- 出力するアクション:
試験ポイント: 運用データパイプラインでは通常、データ書き込みアクションのみを存在させるべき。その他のアクション(
count/display/collectなど)はクエリ最適化を妨げ、ボトルネックの原因になる。
遅延評価(Lazy Evaluation)
- Spark はアクションが呼ばれるまで変換を評価しない。指定順に 1 つずつ評価するのではなく、アクションが全変換の計算をトリガーするまで待つ。
- これにより複数の操作を「チェーン」でき、Spark が最も効率的な**物理プラン(physical plan)**を特定して最適化できる。
- pandas DataFrame の**即時実行(eager execution)**とは大きく異なる。
不変性(Immutability)
- DataFrame は不変。変換は元のデータや DataFrame を変更せず、新しい DataFrame を返す。
- 後続操作でアクセスするには変数に保存する必要がある。
- Scala では特に、DataFrame を変更したら新しい変数に代入しなければならない(
val dfRenamed = df.withColumnRenamed(...))。
狭い変換(narrow)と広い変換(wide)— 性能判断の基礎
| 狭い変換(narrow) | 広い変換(wide) | |
|---|---|---|
| 依存関係 | 各出力パーティションが単一の入力パーティションに依存 | 出力パーティションが複数の入力パーティションに依存 |
| シャッフル | 不要 | 必要(データの再分配) |
| 例 | select, filter, withColumn, map, union | groupBy, join, orderBy, distinct, repartition |
| コスト | 低い | 高い(ネットワーク I/O、ステージ境界を作る) |
試験ポイント: コードを見て「どの操作がシャッフルを起こすか」を判断できること。
groupBy/join/orderByなどの wide 変換が性能ボトルネックになりやすい。
PySpark の API とライブラリ
- Spark SQL と DataFrames: リレーショナルクエリ。SQL と Spark プログラムを混在可能。
- 構造化ストリーミング(Structured Streaming): バッチと同じ書き方でストリーム処理を表現。Spark SQL エンジンが増分的・継続的に実行。
- Spark 上の Pandas API: pandas のコードを分散スケール。テストは pandas、運用は Spark と単一コードベースで対応。
- 機械学習(MLlib): スケーラブルな ML ライブラリ。
- GraphX: グラフ並列計算。
3-2. 主要な DataFrame 操作(select/filter/withColumn/join/groupBy/agg など)
DataFrame の作成
python
# データとスキーマ文字列から作成
data = [[2021, "test", "Albany", "M", 42]]
df1 = spark.createDataFrame(
data,
schema="Year int, First_Name STRING, County STRING, Sex STRING, Count int")
display(df1) # display() は Databricks ノートブック専用(リッチな可視化)
# df1.show() # show() は Apache Spark DataFrame API(基本的な可視化)display()は Databricks ノートブック固有でリッチな可視化を提供。show()は Apache Spark DataFrame API の一部で基本的な可視化。両方ともアクション。
CSV ファイルの読み込み
python
df_csv = spark.read.csv(
f"{path_volume}/{file_name}",
header=True, # 1 行目をヘッダーとして扱う
inferSchema=True, # スキーマを自動推論
sep=",") # 区切り文字
display(df_csv)spark.read.format("json").json("/path")のように多くのファイル形式(CSV / JSON / Parquet など)を読み込める。
スキーマの確認
python
df_csv.printSchema() # 列名とデータ型を表示注: Databricks では「スキーマ」という語を、DataFrame の列定義の意味と、カタログに登録されたテーブルのコレクション(=ネームスペース)の意味の両方で使う。
列の名前変更
python
df_csv = df_csv.withColumnRenamed("First Name", "First_Name")行の追加(union)
python
df = df1.union(df_csv) # 2 つの DataFrame の行を結合
display(df)フィルター(filter / where)
python
# filter() と where() は完全に同義(パフォーマンス・構文とも差なし)
display(df.filter(df["Count"] > 50))
display(df.where(df["Count"] > 50))列の選択と並べ替え(select / orderBy / desc)
python
from pyspark.sql.functions import desc
display(df.select("First_Name", "Count").orderBy(desc("Count")))pyspark.sql.functionsに SQL 関数が用意されている(orderBy/desc/exprなど)。必要に応じてインポートする。
複合条件でのサブセット作成(チェーン)
python
subsetDF = (df
.filter((df["Year"] == 2009) & (df["Count"] > 100) & (df["Sex"] == "F"))
.select("First_Name", "County", "Count")
.orderBy(desc("Count")))
display(subsetDF)- 複数条件は
&(AND)、|(OR)で結合し、各条件を括弧で囲む(演算子優先順位のため必須)。
SQL 式を DataFrame 内で使う
python
# selectExpr(): SQL 式を受け取る select のバリエーション
display(df.selectExpr("Count", "upper(County) as big_name"))
# expr(): 列を指定する任意の場所で SQL 構文を使う
from pyspark.sql.functions import expr
display(df.select("Count", expr("lower(County) as little_name")))
# spark.sql(): 任意の SQL クエリを実行
display(spark.sql(f"SELECT * FROM {path_table}.{table_name}"))DataFrame の保存
python
# テーブルに保存(Databricks は既定で Delta Lake 形式)
df.write.mode("overwrite").saveAsTable(f"{path_table}.{table_name}")
# JSON ファイル(複数ファイルのディレクトリ)に保存
df.write.format("json").mode("overwrite").save("/tmp/json_data")- Spark は分散処理のため、1 ファイルではなくファイルのディレクトリを書き出す。Delta Lake は Parquet フォルダー/ファイルを分割する。
- Databricks はほとんどのアプリでファイルパスよりもテーブルの使用を推奨。
集計・結合の DataFrame API については本節のコード例と 3-3・3-4 を参照。
groupBy(...).agg(...)の具体例は UDAF(3-5)と高階関数(3-4)にも登場する。
3-3. ウィンドウ関数(種類・PARTITION BY/ORDER BY・フレーム指定・具体例)
ウィンドウ関数とは
- ウィンドウと呼ばれる行のグループを操作し、行のグループに基づいて各行の戻り値を計算する関数。
- 集計(
GROUP BY)と違い行数を減らさない。各行に対して 1 つの計算結果を付加する。 - 用途: 現在の行の相対位置を考慮した値へのアクセス、移動平均、累積統計など。
構文
function OVER { window_name | ( window_name ) | window_spec }
function
{ ranking_function | analytic_function | aggregate_function }
window_spec
( [ PARTITION BY partition [ , ... ] ] [ order_by ] [ window_frame ] )window_spec の 3 要素:
- PARTITION BY partition: 関数が作用する行のグループを定義。省略時は全行が 1 パーティション。
- ORDER BY(order_by): パーティション内の行の順序を指定。
- window_frame: パーティション内で集計/分析関数が作用するスライディングサブセットを指定。
エイリアスの注意点:
ORDER BYはSORT BYのエイリアスとして指定可能。PARTITION BYはDISTRIBUTE BYのエイリアスとして指定可能。CLUSTER BYがない場合、ORDER BYをPARTITION BYのエイリアスとして使える。
関数の 3 クラスとフレームの可否
| クラス | 例 | ORDER BY | window_frame |
|---|---|---|---|
| ランキング関数(ranking) | ROW_NUMBER, RANK, DENSE_RANK, NTILE, PERCENT_RANK | 必須 | 指定不可 |
| 分析関数(analytic) | LAG, LEAD, FIRST, LAST, NTH_VALUE, CUME_DIST | 通常必須 | 関数による |
| 集計関数(aggregate) | SUM, AVG, COUNT, MIN, MAX | 任意 | 指定可 |
試験ポイント: ランキング関数に
window_frame(ROWS/RANGE)を付けるとWINDOW_FUNCTION_AND_FRAME_MISMATCHエラー。ランキング関数は ORDER BY を含む window_spec が必須だが、フレーム句は含めてはならない。集計関数にFILTER句は使えない。
ランキング関数の違い(最頻出)
同順位(タイ)の扱いが異なる。
| 関数 | 同順位の扱い | 次の順位 |
|---|---|---|
ROW_NUMBER() | タイでも連番(一意) | 常に +1 |
RANK() | タイは同順位 | タイの数だけ飛ぶ(1,2,2,4) |
DENSE_RANK() | タイは同順位 | 飛ばない(1,2,2,3) |
公式例(部署ごとに給与で順位付け、RANK):
sql
SELECT name,
dept,
RANK() OVER (PARTITION BY dept ORDER BY salary) AS rank
FROM employees;
-- Engineering: Fred 21000→1, Tom 23000→2, Chloe 23000→2, Paul 29000→4
-- ↑ 23000 が 2 人(2 位タイ)なので次は 4 位に飛ぶ同じデータを DENSE_RANK にすると:
sql
SELECT name, dept, salary,
DENSE_RANK() OVER (PARTITION BY dept ORDER BY salary) AS dense_rank
FROM employees;
-- Engineering: Fred→1, Tom→2, Chloe→2, Paul→3
-- ↑ タイの後も飛ばず 3 位分析関数 LAG / LEAD
LAG(col): 現在の行より前の行の値を取得。LEAD(col): 現在の行より後の行の値を取得。LAG(col, offset, default)/LEAD(col, offset, default): オフセットと、範囲外時のデフォルト値を指定可能。
sql
SELECT name, salary,
LAG(salary) OVER (PARTITION BY dept ORDER BY salary) AS lag,
LEAD(salary, 1, 0) OVER (PARTITION BY dept ORDER BY salary) AS lead
FROM employees;
-- Sales: Lisa 10000 → lag=NULL(前がない), lead=30000
-- Alex 30000 → lag=10000, lead=32000
-- Evan 32000 → lag=30000, lead=0(後がないのでデフォルト0)CUME_DIST(累積分布)
sql
SELECT name, dept, age,
CUME_DIST() OVER (PARTITION BY dept ORDER BY age
RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS cume_dist
FROM employees;
-- Engineering (4名): 23→0.25, 25→0.50, 28→0.75, 33→1.0集計ウィンドウ関数(MIN の例)
sql
SELECT name, dept, salary,
MIN(salary) OVER (PARTITION BY dept ORDER BY salary) AS min
FROM employees;
-- 各部署の最小給与が全行に付加されるウィンドウフレーム(window frame)— ROWS と RANGE
集計/分析ウィンドウ関数が作用する、パーティション内の行のスライディングサブセットを指定する。
構文
{ frame_mode frame_start |
frame_mode BETWEEN frame_start AND frame_end }
frame_mode : { RANGE | ROWS }
frame_start : { UNBOUNDED PRECEDING | offset_start PRECEDING | CURRENT ROW | offset_start FOLLOWING }
frame_end : { offset_stop PRECEDING | CURRENT ROW | offset_stop FOLLOWING | UNBOUNDED FOLLOWING }フレーム境界(frame boundary)
| 境界 | 意味 |
|---|---|
UNBOUNDED PRECEDING | パーティションの先頭から |
offset PRECEDING | 現在行の offset 個前から/まで |
CURRENT ROW | 現在の行 |
offset FOLLOWING | 現在行の offset 個後から/まで |
UNBOUNDED FOLLOWING | パーティションの末尾まで |
frame_endを省略するとCURRENT ROWで停止。- 末尾は始点より後でなければならない。
ROWS と RANGE の決定的な違い(最頻出)
- ROWS: 現在の行の前後の物理的な行数でフレームを表す。タイ(同じ ORDER BY 値)の行も個別にカウント。
- RANGE: 現在の行の ORDER BY 値(
obExpr)からの値のオフセットでフレームを表す。同じ ORDER BY 値を持つ行をまとめて扱う。
RANGE の制約:
ORDER BY句が必須、かつ式は 1 つのみ。
- ORDER BY なし →
DATATYPE_MISMATCH.RANGE_FRAME_WITHOUT_ORDER- ORDER BY が複数式 →
DATATYPE_MISMATCH.RANGE_FRAME_MULTI_ORDER
ROWS vs RANGE の挙動差(公式例・Engineering 部門、給与で累積和)
sql
SELECT name, salary,
SUM(salary) OVER (ORDER BY salary
ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS rows_total,
SUM(salary) OVER (ORDER BY salary
RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS range_total
FROM employees WHERE dept = 'Engineering';
-- name salary rows_total range_total
-- Fred 21000 21000 21000
-- Chloe 23000 44000 67000 ← ROWS は 1 行ずつ(21000+23000), RANGE は 23000 のタイをまとめる
-- Tom 23000 67000 67000 ← ROWS は 3 行目まで, RANGE は同値をまとめて同じ値
-- Paul 29000 96000 96000- Chloe と Tom は同じ給与 23000。ROWS では各行が個別カウントされ累積が異なる(44000 vs 67000)。RANGE では同値をまとめるため両者とも同じ 67000 になる。
代表的なフレームパターン
sql
-- 累積和(running total)
SUM(salary) OVER (PARTITION BY dept ORDER BY salary
ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW)
-- 3 行移動平均(前1行 + 現在行 + 後1行)
ROUND(AVG(salary) OVER (PARTITION BY dept ORDER BY salary
ROWS BETWEEN 1 PRECEDING AND 1 FOLLOWING))
-- 現在行から末尾までの合計
SUM(salary) OVER (PARTITION BY dept ORDER BY salary
ROWS BETWEEN CURRENT ROW AND UNBOUNDED FOLLOWING)
-- 値ベースの範囲(給与 ±5000 の合計)
SUM(salary) OVER (PARTITION BY dept ORDER BY salary
RANGE BETWEEN 5000 PRECEDING AND 5000 FOLLOWING)デフォルトフレームの注意(重要): 集計ウィンドウ関数で
ORDER BYを付けフレームを明示しない場合、デフォルトはRANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW(累積計算)になる。ORDER BYを付けない場合はパーティション全体が対象。「全期間の合計を出したいのに ORDER BY を付けたら累積和になってしまう」というのは典型的な落とし穴。
PySpark でのウィンドウ定義
python
from pyspark.sql import Window
from pyspark.sql.functions import row_number, rank, dense_rank, sum as _sum
w = (Window
.partitionBy("dept")
.orderBy("salary")
.rowsBetween(Window.unboundedPreceding, Window.currentRow))
df.withColumn("running_total", _sum("salary").over(w))
# 全パーティションを対象にする(フレーム全体)
w_all = (Window
.partitionBy("id")
.rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing))3-4. 複雑・ネストしたデータの処理(struct/array/map、explode、高階関数)
Spark の列は単純型だけでなく**複雑型(complex / nested types)**を表せる。
複雑型
- struct(構造体): 名前付きフィールドを持つネストしたレコード。
col.fieldまたはcol["field"]でアクセス。DataFrame API ではstruct(...)で作成。 - array(配列): 同じ型の要素の順序付きコレクション。
array(...)で作成、col[0]で要素アクセス、size(col)で長さ。 - map(マップ): キーと値のペア。
map(...)で作成、col[key]で値アクセス。
explode 系関数(ネストの展開)
| 関数 | 挙動 |
|---|---|
explode(array_or_map) | 各要素を個別の行に展開。空/NULL の要素は行が消える |
explode_outer(...) | explode と同じだが、空/NULL でも NULL 行を保持 |
posexplode(...) | 要素と併せて位置(インデックス)も返す |
inline(array_of_structs) | 構造体配列を複数列に展開 |
sql
-- SQL: 配列を行に展開
SELECT id, explode(items) AS item
FROM orders;python
# PySpark: 配列を行に展開
from pyspark.sql.functions import explode, col
df.select("id", explode(col("items")).alias("item"))- 逆操作は
collect_list/collect_set(行を配列に集約する集計関数)。
高階関数(Higher-Order Functions)とラムダ式
- Azure Databricks は Apache Spark SQL で配列を操作する専用プリミティブを提供。定型コードを大量に書かずに配列操作を簡潔に記述できる。
- 2 つの関数型プログラミング構成要素を中心に展開する:
- 高階関数(higher-order function): 配列を受け取り、その処理方法を実装し、計算結果を指示する。
- 匿名(ラムダ)関数(lambda function): 配列内の各項目をどう処理するかを委任される無名関数。
- Apache Spark には、高階関数を含め配列型などの複合型を操作する組み込み関数がある。
代表的な高階関数(Spark SQL):
| 関数 | 説明 | 例 |
|---|---|---|
transform(array, x -> expr) | 各要素に式を適用して新しい配列を返す | transform(vals, x -> x + 1) |
filter(array, x -> predicate) | 述語が真の要素だけ残す | filter(vals, x -> x > 0) |
exists(array, x -> predicate) | いずれかが述語を満たせば true | exists(vals, x -> x > 100) |
forall(array, x -> predicate) | すべてが述語を満たせば true | forall(vals, x -> x > 0) |
aggregate(array, start, (acc, x) -> ..., acc -> ...) | 畳み込みで集約 | aggregate(vals, 0, (acc, x) -> acc + x) |
zip_with(a1, a2, (x, y) -> expr) | 2 配列を要素ごとに結合 | zip_with(a, b, (x, y) -> x + y) |
sql
-- 各要素を 2 乗し、正の値だけ残して合計する
SELECT
aggregate(
filter(transform(values, x -> x * x), y -> y > 0),
0,
(acc, v) -> acc + v
) AS result
FROM data;試験ポイント: 高階関数は
explode→ 集計 →collect_listのようなシャッフルを伴う往復を避けて配列を直接操作できるため、ネストデータ処理で性能面のメリットがある。ラムダ式の構文x -> expr、2 引数(acc, x) -> exprを読めるようにしておく。
JSON / 半構造化データ
from_json(col, schema)で JSON 文字列を struct に、to_json(col)で struct を JSON 文字列に変換。get_json_object/json_tupleで JSON 文字列から値を抽出。- Databricks には VARIANT 型(半構造化データ用)もあり、UDF から
VariantVal.parseJson(...)で生成できる(3-5 のバリアント例参照)。
3-5. UDF の種類と性能特性(スカラー/pandas/UDTF、Unity Catalog 管理、性能順)
UDF を使う場面
- UDF を使う: 組み込みの Apache Spark 関数では表現が困難なロジック。アドホッククエリ、手動データクリーニング、探索的データ分析、小〜中規模データセット。一般的用途は暗号化・復号・ハッシュ・JSON 解析・検証。
- 組み込み関数 / Spark メソッドを使う: 大規模データセット、ETL ジョブ、ストリーミングなど定期的・継続的に実行されるワークロード。組み込み関数は分散処理向けに最適化されており大規模で高速。
UDF の種類
| 種類 | 入出力 | スコープ | 特徴 |
|---|---|---|---|
| スカラー UDF(Scalar) | 1 行 → 1 値 | UC 管理 or セッション | 各行に 1 結果 |
| バッチスカラー UDF(Batch Scalar) | バッチ(1:1 行パリティ維持) | UC 管理 or セッション | 行ごとオーバーヘッド軽減、バッチ間で状態保持可 |
| 非スカラー UDF / pandas UDF | 1:N・多:多 | セッション | Arrow でベクトル化、高速 |
| UDAF | 複数行 → 1 集計値 | セッション限定 | 集計 |
| UDTF | 1 行 → 複数行(複数列) | UC 管理 or セッション | テーブルを返す |
スカラー UDF(SQL / Python)の例
sql
-- SQL UDF(Unity Catalog に永続化)
CREATE OR REPLACE FUNCTION main.test.get_name_length(name STRING)
RETURNS INT
RETURN LENGTH(name);
SELECT name, main.test.get_name_length(name) AS name_length FROM your_table;python
# PySpark スカラー UDF(デコレータ構文)
from pyspark.sql.functions import udf
from pyspark.sql.types import IntegerType
@udf(returnType=IntegerType())
def get_name_length(name):
return len(name)
df = df.withColumn("name_length", get_name_length(df.name))Python UDF の登録と SQL からの呼び出し
python
def squared(s):
return s * s
# SQL から使えるよう登録(既定の戻り値型は StringType)
spark.udf.register("squaredWithPython", squared)
# 戻り値型を明示
from pyspark.sql.types import LongType
spark.udf.register("squaredWithPython", squared, LongType())sql
%sql select id, squaredWithPython(id) as id_squared from testpython
# DataFrame 内で使う(udf() またはアノテーション @udf("long"))
from pyspark.sql.functions import udf
@udf("long")
def squared_udf(s):
return s * s
display(df.select("id", squared_udf("id").alias("id_squared")))バッチ Unity Catalog Python UDF(PARAMETER STYLE PANDAS)
sql
CREATE OR REPLACE FUNCTION main.test.calculate_bmi_pandas(weight_kg DOUBLE, height_m DOUBLE)
RETURNS DOUBLE
LANGUAGE PYTHON
PARAMETER STYLE PANDAS
HANDLER 'handler_function'
AS $$
import pandas as pd
from typing import Iterator, Tuple
def handler_function(batch_iter: Iterator[Tuple[pd.Series, pd.Series]]) -> Iterator[pd.Series]:
for weight_series, height_series in batch_iter:
yield weight_series / (height_series ** 2)
$$;UDTF(テーブルを返す)
sql
CREATE OR REPLACE FUNCTION get_sum_diff(x INT, y INT)
RETURNS TABLE (sum INT, diff INT)
LANGUAGE PYTHON
HANDLER 'GetSumDiff'
AS $$
class GetSumDiff:
def eval(self, x: int, y: int):
yield x + y, x - y
$$;
SELECT * FROM get_sum_diff(10, 3); -- → sum=13, diff=7python
# PySpark UDTF
from pyspark.sql.functions import lit, udtf
@udtf(returnType="sum: int, diff: int")
class GetSumDiff:
def eval(self, x: int, y: int):
yield x + y, x - y
GetSumDiff(lit(1), lit(2)).show()UDAF(集計・セッションスコープ)
python
from pyspark.sql.functions import pandas_udf
import pandas as pd
@pandas_udf("int")
def total_score_udf(scores: pd.Series) -> int:
return scores.sum()
result_df = (df.groupBy("name_length")
.agg(total_score_udf(df["score"]).alias("total_score")))pandas UDF(ベクトル化 UDF)— 種類と型ヒント
- Apache Arrow でデータ転送し pandas で処理。1 行ずつの Python UDF に比べ最大 100 倍高速になり得るベクトル化操作が可能。
@pandas_udfデコレータと Python 型ヒントで種類を指定する。
| 種類 | 型ヒント | 用途 |
|---|---|---|
| Series → Series | pd.Series -> pd.Series | select / withColumn でのスカラー操作。入力と同じ長さの Series を返す |
| Iterator[Series] → Iterator[Series] | Iterator[pd.Series] -> Iterator[pd.Series] | 1 列入力。状態初期化(例: ML モデルのロード)に有用 |
| Iterator[複数 Series] → Iterator[Series] | Iterator[Tuple[pd.Series, ...]] -> Iterator[pd.Series] | 複数列入力 |
| Series → Scalar | pd.Series, ... -> Any | 集計。groupBy.agg / Window で使用。部分集計非対応、各グループ全体をメモリにロード |
python
# Series to Series
import pandas as pd
from pyspark.sql.functions import pandas_udf
@pandas_udf("double")
def calculate_bmi_pandas(weight: pd.Series, height: pd.Series) -> pd.Series:
return weight / (height ** 2)
df.withColumn("BMI", calculate_bmi_pandas(df["Weight"], df["Height"])).display()python
# Iterator of Series to Iterator of Series(状態を初期化できる)
from typing import Iterator
@pandas_udf("long")
def plus_one(batch_iter: Iterator[pd.Series]) -> Iterator[pd.Series]:
for x in batch_iter:
yield x + 1python
# Series to Scalar(集計)
from pyspark.sql import Window
@pandas_udf("double")
def mean_udf(v: pd.Series) -> float:
return v.mean()
df.select(mean_udf(df["v"])).show() # 全体平均
df.groupby("id").agg(mean_udf(df["v"])).show() # グループ平均
w = Window.partitionBy("id").rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing)
df.withColumn("mean_v", mean_udf(df["v"]).over(w)).show() # ウィンドウ平均- 関連 API:
mapInPandas、mapInArrow、applyInPandas(Grouped Map)。非スカラー UDF に分類される。 - Arrow バッチサイズ:
spark.sql.execution.arrow.maxRecordsPerBatch(既定 1 万レコード/バッチ)で調整可能。列数が多い場合は下げる。サーバーレスや DBR 13.3〜14.2 標準アクセスでは無効(プラットフォームが内部管理)。 - タイムゾーン: Spark は内部で UTC 保持。pandas は
datetime64[ns]。toPandas()や pandas UDF で自動変換される。
Unity Catalog 管理 UDF vs セッションスコープ UDF
| 考慮事項 | Unity Catalog 管理 UDF | セッションスコープ UDF |
|---|---|---|
| 最適な用途 | チーム・ノートブック・ジョブ・SQL ウェアハウス間で安全に共有 | 1 ノートブック/ジョブ内での迅速な反復開発 |
| 対応言語 | SQL, Python, Scala, Java | SQL, Python, Scala |
| ガバナンス・共有 | UC の権限で管理、カタログエクスプローラーで検出可能 | 現在の SparkSession 限定。管理も共有もされない |
| 永続性 | UC に永続化、セッション間で再利用可 | 現在のセッションのみ |
性能に関する考慮事項(最頻出・性能順)
効率のよい順(速い → 遅い):組み込み関数(built-in functions) = SQL UDF > Scala UDF > pandas UDF > Python UDF(行ごと)
- 組み込み関数と SQL UDF が最も効率的。可能な限りこれらを使う。
- Scala UDF は通常 Python UDF より高速。
- 非分離の Scala UDF は JVM 上で実行され、JVM 内外へのデータ移動オーバーヘッドを回避。
- 分離された Scala UDF は JVM ⇔ データ移動が必要だが、メモリ処理が効率的で Python UDF より高速。
- Python UDF / pandas UDF は Scala UDF より遅い傾向。理由: データをシリアライズして JVM から Python インタプリタに移動する必要があるため。
- pandas UDF は Python UDF より最大 100 倍高速。理由: Apache Arrow でシリアライズコストを削減するベクトル化。
UDF の注意点・落とし穴
- 評価順序と NULL チェック(重要): Spark SQL では部分式の評価順序が保証されない。
AND/ORに左から右の短絡(ショートサーキット)セマンティクスはない。WHERE s IS NOT NULL AND strlen(s) > 1でも、NULL 除外後に UDF が呼ばれる保証はない。- 対策 1: UDF 自体を NULL 対応にする(内部で NULL チェック)。
- 対策 2:
IF/CASE WHENで NULL チェックし、条件分岐内で UDF を呼ぶ。
python
spark.udf.register("strlen_nullsafe", lambda s: len(s) if s is not None else -1, "int")
spark.sql("select s from test1 where s is not null and strlen_nullsafe(s) > 1") # OK
spark.sql("select s from test1 where if(s is not null, strlen(s), null) > 1") # OK- メモリ制限: サーバーレスの PySpark UDF は 1 GB/UDF(超過で
UDF_PYSPARK_USER_CODE_ERROR.MEMORY_LIMIT_SERVERLESS)。標準アクセスモードはインスタンスの利用可能メモリに依存。 - ブロードキャスト変数: 標準アクセスモードクラスター・サーバーレスの PySpark UDF は非対応。
- ネットワークアクセス: サーバーレス SQL ウェアハウスの Python UDF は既定で送信ネットワーク要求不可(ハングする)。
- ファイルアクセス: DBR 14.2 以下の共有クラスターの PySpark UDF は Git フォルダー / ワークスペースファイル / UC ボリュームにアクセス不可。
- ランタイム要件: スカラー Python UDF / pandas UDF は DBR 13.3 LTS 以降で全アクセスモード対応。DBR 12.2 LTS 以下では UC 標準アクセスモードで非対応。
4. 構文・コード例
4-1. DataFrame チェーンの基本形
python
from pyspark.sql.functions import col, desc
result = (spark.read
.csv("/Volumes/cat/sch/vol/data.csv", header=True, inferSchema=True) # 変換
.filter(col("Count") > 50) # 変換(narrow)
.select("First_Name", "County", "Count") # 変換(narrow)
.orderBy(desc("Count"))) # 変換(wide=シャッフル)
result.write.mode("overwrite").saveAsTable("cat.sch.result") # アクション(書き込み)4-2. groupBy と集計
python
from pyspark.sql.functions import sum, avg, count, max, min
(df.groupBy("dept")
.agg(sum("salary").alias("total"),
avg("salary").alias("avg"),
count("*").alias("n"),
max("salary").alias("max"))
.display())4-3. join(結合)
python
# 内部結合 / 左外部結合 / 完全外部結合など how で指定
joined = df1.join(df2, on="key", how="inner")
joined = df1.join(df2, df1["k"] == df2["k"], how="left")
# how: inner / left / right / full / left_semi / left_anti / cross4-4. ウィンドウ関数(SQL)
sql
-- ランキング + 累積和 + 移動平均を 1 クエリで
SELECT name, dept, salary,
ROW_NUMBER() OVER (PARTITION BY dept ORDER BY salary DESC) AS rn,
RANK() OVER (PARTITION BY dept ORDER BY salary DESC) AS rnk,
DENSE_RANK() OVER (PARTITION BY dept ORDER BY salary DESC) AS drnk,
SUM(salary) OVER (PARTITION BY dept ORDER BY salary
ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS running_total,
AVG(salary) OVER (PARTITION BY dept ORDER BY salary
ROWS BETWEEN 1 PRECEDING AND 1 FOLLOWING) AS moving_avg,
LAG(salary) OVER (PARTITION BY dept ORDER BY salary) AS prev_salary
FROM employees;4-5. ウィンドウ関数(PySpark)
python
from pyspark.sql import Window
from pyspark.sql.functions import row_number, rank, dense_rank, lag, sum as _sum
w = Window.partitionBy("dept").orderBy(col("salary").desc())
w_run = Window.partitionBy("dept").orderBy("salary") \
.rowsBetween(Window.unboundedPreceding, Window.currentRow)
(df.withColumn("rn", row_number().over(w))
.withColumn("running_total", _sum("salary").over(w_run))
.withColumn("prev", lag("salary").over(Window.partitionBy("dept").orderBy("salary")))
.display())4-6. ネストデータの展開と高階関数
sql
-- 配列を行に展開
SELECT order_id, explode(items) AS item FROM orders;
-- struct フィールドへのアクセス
SELECT customer.name, customer.address.city FROM t;
-- 高階関数: 各要素を2乗し正の値だけ合計
SELECT id,
aggregate(filter(transform(vals, x -> x * x), y -> y > 0), 0, (acc, v) -> acc + v) AS s
FROM t;4-7. 言語間データ共有(一時ビュー)
python
# Python DataFrame → 一時ビュー化
df.createOrReplaceTempView("my_view")sql
%sql
-- SQL セルから同じデータへアクセス
SELECT * FROM my_view WHERE Count > 100;4-8. Python に SQL を埋め込む
python
result = spark.sql(f"SELECT * FROM {catalog}.{schema}.{table} WHERE Year = 2009")
display(result)5. 試験で問われるポイント(コード読解・性能判断重視)
変換 vs アクションの識別: 与えられたコード列のうちどれがジョブをトリガーするか。
select/filter/joinは変換、count/show/collect/write/takeはアクション。アクションを呼ぶまで何も実行されない(遅延評価)。遅延評価の帰結: 「変換だけ書いてアクションを呼ばないと計算は走らない」「運用パイプラインは書き込みアクションのみにすべき」「不要な
count/displayは最適化を妨げる」。narrow vs wide とシャッフル:
groupBy/join/orderBy/distinct/repartitionはシャッフル(wide)を起こし高コスト。select/filter/withColumnは narrow で低コスト。「どの操作が性能ボトルネックか」を問われる。不変性: 「
df.withColumnRenamed(...)を実行しても元のdfは変わらない」「結果を変数に代入しないと反映されない」。Scala では新変数への代入が必須。ウィンドウ関数の結果を計算させる:
ROW_NUMBER/RANK/DENSE_RANKのタイの扱いの差(1,2,3 / 1,2,2,4 / 1,2,2,3)。LAG/LEADの境界(前/後がない行は NULL またはデフォルト値)。- ROWS と RANGE の差(タイ行の扱い。RANGE は同値をまとめる)。
- ORDER BY を付けた集計ウィンドウのデフォルトは累積和(
RANGE UNBOUNDED PRECEDING AND CURRENT ROW)。 - ランキング関数にフレームを付けるとエラー、RANGE フレームにORDER BY なし/複数式でエラー。
UDF の性能順(頻出): 組み込み/SQL UDF > Scala UDF > pandas UDF > Python UDF。「大規模 ETL でどの実装を選ぶべきか」→ まず組み込み関数、無理なら pandas UDF(Arrow ベクトル化)、行ごとの Python UDF は最終手段。
pandas UDF がなぜ速いか: Apache Arrow による列指向転送でシリアライズコストを削減し、ベクトル化するから。「Python UDF の 100 倍高速」。
UDF の種類判別: 入出力から判断。1 行→1 値=スカラー、複数行→1 値=UDAF、1 行→複数行=UDTF、Series→Series/Scalar=pandas UDF。
UDF の NULL 落とし穴:
WHERE x IS NOT NULL AND udf(x)は短絡保証がない → UDF 内で NULL 対応するかIF/CASE WHENを使う。Unity Catalog 管理 vs セッションスコープ: チーム共有・永続化・ガバナンス → UC 管理。一時的な反復開発 → セッションスコープ。UC 管理は Java も可、セッションは SQL/Python/Scala。
言語選択: 新規プロジェクトは Python と SQL 推奨。Scala/R は制限あり非推奨(R はノートブックのみ完全サポート)。Lakeflow パイプラインは Python/SQL、ワークフローは Python/SQL/Scala/Java 対応。統一エンジンのため言語間で性能はほぼ同じ。
ネストデータ処理:
explodeで配列→行、高階関数(transform/filter/aggregate)でシャッフルを避けて配列を直接操作。ラムダ構文x -> expr、(acc, x) -> exprを読めること。
6. 理解度チェックリスト
- [ ] DataFrame が RDD 上に構築された抽象化であり、統一エンジンにより言語間で性能差がほぼないことを説明できる
- [ ] 変換(transformation)とアクション(action)を具体例で分類でき、遅延評価の意味を説明できる
- [ ] 「運用パイプラインは書き込みアクションのみにすべき」理由(他のアクションが最適化を妨げる)を言える
- [ ] DataFrame が不変であること、変換結果を変数に保存する必要があることを理解している
- [ ] narrow 変換と wide 変換を分類し、どれがシャッフルを起こすか判断できる
- [ ]
select/filter/where/withColumn/withColumnRenamed/union/orderBy/selectExpr/expr/spark.sqlの使い方を書ける - [ ]
filter()とwhere()が完全に同義であることを知っている - [ ] 複合条件
(cond1) & (cond2)の括弧が必須である理由を理解している - [ ] ウィンドウ関数の
PARTITION BY/ORDER BY/window_frameの役割を説明できる - [ ]
ROW_NUMBER/RANK/DENSE_RANKのタイ時の出力差を暗算できる - [ ]
LAG/LEADの境界(NULL/デフォルト値)の挙動を答えられる - [ ]
ROWSフレームとRANGEフレームの違い(物理行数 vs 値オフセット、タイ行の扱い)を説明できる - [ ] RANGE フレームは ORDER BY 式が 1 つ必須で、なければエラーになることを知っている
- [ ] ORDER BY 付き集計ウィンドウのデフォルトフレームが累積和になることを知っている
- [ ] ランキング関数にフレーム句を付けるとエラーになることを知っている
- [ ] struct / array / map の複雑型と
explode/explode_outer/posexplodeの違いを説明できる - [ ] 高階関数
transform/filter/aggregate/existsとラムダ構文を読める - [ ] スカラー UDF / バッチスカラー UDF / pandas UDF / UDTF / UDAF を入出力から判別できる
- [ ] UDF の性能順(組み込み/SQL UDF > Scala UDF > pandas UDF > Python UDF)を暗記している
- [ ] pandas UDF が Apache Arrow でベクトル化し Python UDF の最大 100 倍速い理由を説明できる
- [ ] pandas UDF の 4 種類(Series→Series、Iterator、複数 Series、Series→Scalar)を型ヒントで区別できる
- [ ]
@udf/@pandas_udf/spark.udf.register/ SQL のCREATE FUNCTIONの書き方を知っている - [ ] UDF での NULL チェックの落とし穴(評価順序が保証されない)と対策を説明できる
- [ ] Unity Catalog 管理 UDF とセッションスコープ UDF の違い(共有・永続化・言語)を説明できる
- [ ]
createOrReplaceTempViewによる言語間データ共有とspark.sqlの埋め込みを使える - [ ] 新規プロジェクトで Python / SQL が推奨され、Scala / R が制限付きである理由を言える