Skip to content

Professional 学習教材 ① Python/SQL でのデータ処理コード開発(配点 22%)

Databricks 公式ドキュメント(日本語版)の内容を、重要用語・概念を漏らさずまとめた Professional 試験向け自習教材です。 参照した公式ページ:

注: 高階関数ページ(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 / フレーム指定 ROWSRANGE)による累積・移動平均・ランキング計算
  • 複雑・ネストしたデータstruct / array / mapexplode、高階関数)の処理
  • 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 SparkApache Sparkビッグデータ・機械学習向けの統合分析エンジン。Azure Databricks の基盤。
PySparkPySparkPython から Apache Spark を操作する API。Python と Spark の機能を組み合わせたもの。
データフレームDataFrame名前付き列に編成された 2 次元のラベル付きデータ構造。スプレッドシート/SQL テーブルに相当。Spark の主オブジェクト。
RDD(回復力のある分散データセット)Resilient Distributed Dataset (RDD)DataFrame の基盤となる低レベルの分散データ抽象化。DataFrame は RDD 上に構築される。
スキーマSchemaDataFrame の列の名前と型の定義。手動定義、データソースからの読み取り、推論(inferSchema)が可能。
RowDataFrame 内のレコード。Row オブジェクトとして表される。最適化のためキャッシュとシャッフルに使われる。
Columnスプレッドシートの列に相当。単純型(文字列・整数)だけでなく配列・マップ・null などの複雑型も表す。
変換Transformationデータを読み込み・操作する手順(spark.read / join / 集計 / 型キャストなど)。遅延評価され、新しい DataFrame を返す。
アクションAction変換の結果の計算をトリガーし値を返す操作(display / show / take / count / saveAsTable など)。
遅延評価(遅延読み込み)Lazy Evaluation / Lazy Loading変換はアクションが呼ばれるまで実行されない。論理クエリを一連の命令として保持し、最も効率的な物理プランを特定してから実行。
即時実行Eager Executionpandas DataFrame のモデル。定義した時点で即座に実行される。Spark の遅延評価とは対照的。
不変性ImmutabilityDataFrame は不変。変換は元を変更せず新しい DataFrame を返すため、結果は変数に保存する必要がある。
実行計画 / 物理プランExecution Plan / Physical PlanSpark が変換ロジックを評価するために特定する最も効率的な実行手順。Catalyst オプティマイザが生成。
狭い変換Narrow Transformation各出力パーティションが単一の入力パーティションのみに依存(select / filter / withColumn など)。シャッフル不要で高速。
広い変換Wide Transformation出力パーティションが複数の入力パーティションに依存(groupBy / join / orderBy / distinct など)。シャッフルを伴う。
シャッフルShuffleパーティション間でのデータの再分配。ネットワーク I/O が発生する高コスト操作。ステージ境界を作る。
Spark SQLSpark SQLSQL クエリと Spark プログラムを混在できるリレーショナル処理モジュール。
SparkSessionSparkSessionSpark 機能へのエントリポイント。ノートブックでは 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 FunctionROW_NUMBER / RANK / DENSE_RANK / NTILE / PERCENT_RANK。ORDER BY 必須、フレーム指定不可。
分析関数Analytic FunctionLAG / LEAD / FIRST / LAST / NTH_VALUE / CUME_DIST。前後の行の値を参照。
集計関数Aggregate FunctionSUM / AVG / COUNT / MIN / MAX。ウィンドウと組み合わせて累積・移動計算に使う。
構造体struct名前付きフィールドを持つネストしたレコード型。col.field でアクセス。
配列array同じ型の要素の順序付きコレクション。explode で行に展開可能。
マップmapキーと値のペアのコレクション。
explodeexplode配列/マップの各要素を個別の行に展開する関数。
高階関数Higher-Order Function配列を受け取り、ラムダ関数を各要素に適用する関数(transform / filter / aggregate / exists など)。
ラムダ関数(匿名関数)Lambda / Anonymous Function高階関数に渡す、要素ごとの処理を定義する無名関数。
ユーザー定義関数User-Defined Function (UDF)組み込み関数で表現困難なカスタムロジックを再利用・共有するための関数。
スカラー UDFScalar UDF1 行を処理し、各行に 1 つの結果値を返す UDF。
バッチスカラー UDFBatch Scalar UDF1:1 の入出力を維持しつつ行をバッチ処理する UDF。行ごとのオーバーヘッドを軽減。
pandas UDF(ベクトル化 UDF)pandas UDF / Vectorized UDFApache Arrow でデータ転送し pandas で処理。Python UDF より最大 100 倍高速。
UDTF(ユーザー定義テーブル関数)User-Defined Table Function (UDTF)入力行ごとに複数の行(複数の列)を返す UDF。
UDAF(ユーザー定義集計関数)User-Defined Aggregate Function (UDAF)複数行を処理し 1 つの集計結果を返す UDF。セッションスコープ限定。
Apache ArrowApache Arrow言語をまたぐ列指向インメモリ形式。pandas UDF でシリアライズコストを削減。
Unity Catalog 管理 UDFUnity Catalog-governed UDFUC に永続化され、権限管理・共有・検出が可能な UDF。
セッションスコープ UDFSession-scoped UDF現在の SparkSession に限定される一時的な UDF。
シリアライズのオーバーヘッドSerialization Overheadデータを JVM ⇔ Python インタプリタ間で移動する際のコスト。Python UDF が遅い主因。
一時ビューTemporary View / Temp ViewcreateOrReplaceTempView で作成する、言語間でデータを共有できるセッションスコープのビュー。
マジックコマンド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 などの複雑型も表せる。.dropselect での省略で結果から除かれるだけで、データソースから列が物理的に削除されることはない

変換(Transformation)とアクション(Action)— Professional 最重要概念

Spark は**遅延評価(lazy evaluation)**で変換とアクションを処理する。

変換(Transformation)

  • 処理ロジックを表現する手順。DataFrame を返す。
  • 一般的な変換: データ読み取り(spark.read / spark.table)、結合、集計、型キャスト。
  • すぐには実行されない。論理クエリを「データソースに対する一連の命令」として格納する(メモリ内の結果としてではない)。

アクション(Action)

  • 1 つ以上の DataFrame で一連の変換の結果を計算するよう Spark に指示する。値を返す。
  • 種類:
    • 出力するアクション: displayshow
    • データを収集するアクション(Row を返す): take(n)firsthead
    • データソースに書き込むアクション: 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, uniongroupBy, 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 BYSORT BY のエイリアスとして指定可能。
  • PARTITION BYDISTRIBUTE BY のエイリアスとして指定可能。
  • CLUSTER BY がない場合、ORDER BYPARTITION BY のエイリアスとして使える。

関数の 3 クラスとフレームの可否

クラスORDER BYwindow_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_frameROWS/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)いずれかが述語を満たせば trueexists(vals, x -> x > 100)
forall(array, x -> predicate)すべてが述語を満たせば trueforall(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 UDF1:N・多:多セッションArrow でベクトル化、高速
UDAF複数行 → 1 集計値セッション限定集計
UDTF1 行 → 複数行(複数列)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 test
python
# 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=7
python
# 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 → Seriespd.Series -> pd.Seriesselect / 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 → Scalarpd.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 + 1
python
# 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: mapInPandasmapInArrowapplyInPandas(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, JavaSQL, 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 / cross

4-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. 試験で問われるポイント(コード読解・性能判断重視)

  1. 変換 vs アクションの識別: 与えられたコード列のうちどれがジョブをトリガーするか。select/filter/join は変換、count/show/collect/write/take はアクション。アクションを呼ぶまで何も実行されない(遅延評価)。

  2. 遅延評価の帰結: 「変換だけ書いてアクションを呼ばないと計算は走らない」「運用パイプラインは書き込みアクションのみにすべき」「不要な count/display は最適化を妨げる」。

  3. narrow vs wide とシャッフル: groupBy/join/orderBy/distinct/repartition はシャッフル(wide)を起こし高コスト。select/filter/withColumn は narrow で低コスト。「どの操作が性能ボトルネックか」を問われる。

  4. 不変性: 「df.withColumnRenamed(...) を実行しても元の df は変わらない」「結果を変数に代入しないと反映されない」。Scala では新変数への代入が必須。

  5. ウィンドウ関数の結果を計算させる:

    • 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 なし/複数式でエラー
  6. UDF の性能順(頻出): 組み込み/SQL UDF > Scala UDF > pandas UDF > Python UDF。「大規模 ETL でどの実装を選ぶべきか」→ まず組み込み関数、無理なら pandas UDF(Arrow ベクトル化)、行ごとの Python UDF は最終手段。

  7. pandas UDF がなぜ速いか: Apache Arrow による列指向転送でシリアライズコストを削減し、ベクトル化するから。「Python UDF の 100 倍高速」。

  8. UDF の種類判別: 入出力から判断。1 行→1 値=スカラー、複数行→1 値=UDAF、1 行→複数行=UDTF、Series→Series/Scalar=pandas UDF。

  9. UDF の NULL 落とし穴: WHERE x IS NOT NULL AND udf(x) は短絡保証がない → UDF 内で NULL 対応するか IF/CASE WHEN を使う。

  10. Unity Catalog 管理 vs セッションスコープ: チーム共有・永続化・ガバナンス → UC 管理。一時的な反復開発 → セッションスコープ。UC 管理は Java も可、セッションは SQL/Python/Scala。

  11. 言語選択: 新規プロジェクトは Python と SQL 推奨。Scala/R は制限あり非推奨(R はノートブックのみ完全サポート)。Lakeflow パイプラインは Python/SQL、ワークフローは Python/SQL/Scala/Java 対応。統一エンジンのため言語間で性能はほぼ同じ。

  12. ネストデータ処理: 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 が制限付きである理由を言える