Skip to content

Associate 学習教材 ③ データ変換・モデリング

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


1. このセクションの概要

このセクションは Databricks 認定データエンジニア Associate 試験の最重要領域である「データ変換・モデリング」を扱います。中心となるのは次の 6 テーマです。

  1. Apache Spark を使った ETL(抽出・変換・読み込み / Extract, Transform, Load) の基本。データを DataFrame として読み込み、変換し、Delta Lake に書き込む一連の流れ。
  2. PySpark DataFrame API の主要操作(select / filter / withColumn / join / groupBy など)と、変換(transformation)とアクション(action)/遅延評価(lazy evaluation) の仕組み。
  3. Spark SQL / SQL 言語 によるテーブル・ビューの操作。マネージドテーブル(managed table)/外部テーブル(external table)/ビュー(view)/CTAS(CREATE TABLE AS SELECT) の使い分け。
  4. UDF(ユーザー定義関数 / user-defined function) の種類(SQL UDF、Python UDF、pandas UDF、Scala UDF など)と使いどころ・性能特性。
  5. MERGE INTO によるアップサート(upsert)WHEN MATCHED / WHEN NOT MATCHED / WHEN NOT MATCHED BY SOURCE の 3 句と、SCD Type 2・重複排除・CDC 適用への応用。
  6. メダリオンアーキテクチャ(medallion architecture)。Bronze / Silver / Gold の各層の役割と設計指針。

これらは互いに密接に関連しています。たとえば「Auto Loader で Bronze に取り込み → DataFrame API / SQL で Silver に洗浄・検証 → MERGE INTO や集計で Gold に整形」という流れがメダリオンアーキテクチャの実装そのものになります。

前提となる基盤: Databricks で作成されるテーブルの既定フォーマットは Delta Lake です。Delta Lake は ACID トランザクション(原子性・一貫性・分離性・持続性)を提供するオープンソースのストレージレイヤーであり、データレイクハウス(Data Lakehouse)を実現します。MERGE INTO などの操作は Delta テーブルでのみサポートされます。


2. 重要用語集

用語(日本語)English説明
データフレームDataFrame名前付き列を持つ分散データセット。表形式のデータを表現する Spark の中心的な抽象。spark.read.table() などで取得し、display(df) でプレビューする。
パイスパークPySparkApache Spark の Python API。DataFrame を Python から操作する。
スパーク SQLSpark SQLSpark 上で SQL を実行する仕組み。DataFrame 操作と等価な処理を SQL で記述できる。
ETLExtract, Transform, Load抽出・変換・読み込み。データソースからデータを取り込み、変換し、ターゲット(Delta テーブル等)へ書き込む処理。
変換transformationDataFrame から新しい DataFrame を生成する操作(select, filter, withColumn, join, groupBy など)。遅延評価され、即座には実行されない。
アクションaction実際に計算を起動して結果を返す/書き出す操作(display, count, collect, write/writeStream, toTable など)。
遅延評価lazy evaluation変換は定義されるだけで即実行されず、アクションが呼ばれたときにまとめて最適化・実行される Spark の実行モデル。
自動ローダーAuto Loaderクラウドオブジェクトストレージに到着した新しいファイルを自動検出・処理する増分インジェスト機能。format("cloudFiles") を使う。
デルタレイクDelta LakeACID トランザクションを提供するオープンソースのストレージレイヤー。Databricks の既定テーブル形式。
チェックポイントcheckpoint / checkpointLocation構造化ストリーミングで処理済み位置・状態を記録する場所。増分処理の再開に使う。
構造化ストリーミングStructured StreamingreadStream / writeStream を用いたストリーム処理エンジン。増分処理を担う。
ユーザー定義関数UDF (User-Defined Function)組み込み関数で表現しにくいロジックを再利用可能にするカスタム関数。
スカラー UDFscalar UDF1 行を入力し、各行に 1 つの結果値を返す UDF。
pandas UDFpandas UDFApache Arrow でデータを転送しベクトル化処理する UDF。行ごとの Python UDF より最大 100 倍高速。
UDAFUser-Defined Aggregate Function複数行を入力し 1 つの集計結果を返す UDF。
UDTFUser-Defined Table Function1 つ以上の引数を受け取り、入力行ごとに複数行(複数列も可)を返す UDF。
マージMERGE INTOソーステーブルに基づき、ターゲット Delta テーブルへ更新・挿入・削除をまとめて適用する文。
アップサートupsertUPDATE と INSERT を組み合わせた操作(一致すれば更新、なければ挿入)。MERGE INTO で実現。
メダリオンアーキテクチャmedallion architectureBronze → Silver → Gold と段階的にデータ品質を高めるデータ設計パターン。マルチホップアーキテクチャとも呼ぶ。
ブロンズ層Bronze layer生(未加工・未検証)データを元の形式で取り込む層。
シルバー層Silver layerクレンジング・検証・重複排除・結合を行い、消費しやすい形に整えた層。
ゴールド層Gold layer集計・ディメンションモデリングを施した、分析/BI/ML 向けの高度に洗練された層。
マネージドテーブルmanaged tableデータとメタデータの両方を Databricks が管理するテーブル。DROP でデータ本体も削除される。
外部テーブルexternal table / unmanaged tableデータを外部の既存ストレージ場所に置くテーブル。DROP してもデータ本体は残る。MERGE のターゲットにはできない。
ビューview物理データを持たない仮想テーブル。クエリの結果セットに名前を付けたもの。
一時ビューtemporary view作成したセッション内でのみ有効で、セッション終了時に消えるビュー。
グローバル一時ビューglobal temporary viewglobal_temp スキーマに関連付けられる一時ビュー。
CTASCREATE TABLE AS SELECTクエリ結果から新しいテーブルを作成する構文。
ディメンションモデリングdimensional modelingファクト/ディメンションでビジネスを表現するデータモデリング手法。Gold 層で用いる。
マテリアライズドビューmaterialized view集計結果などを事前計算して保持するビュー。Gold 層の集計に利用。
SCD Type 2Slowly Changing Dimension Type 2変更履歴を新しい行として保持するディメンション管理手法。MERGE INTO で実装。
CDCChange Data Capture変更データ(挿入・更新・削除)を捕捉して反映する処理。MERGE の代表的ユースケース。
スキーマ進化schema evolutionターゲットのスキーマをソースに合わせて自動更新する仕組み(WITH SCHEMA EVOLUTION)。

3. 詳細解説

3-1. Spark を使った ETL の基本(DataFrame の読み書き・変換・アクション・遅延評価)

Databricks 上での ETL は、Apache Spark を使ってデータをオーケストレーションする処理です。公式チュートリアルの流れは次のとおりです。

  1. コンピューティングリソースを作成する(サーバーレスコンピューティングが有効なら不要。ノートブック実行時に自動アタッチされる)。
  2. ノートブックを作成する(ロジックはセル単位で実行。Shift + Enter でセルを実行)。
  3. Auto Loader を構成して Delta Lake へ増分インジェストする。
  4. データを処理・操作する。
  5. ノートブックをジョブとしてスケジュールする(本番スクリプト化)。

インジェスト(抽出): Auto Loader

Databricks では、増分データインジェスト(incremental data ingestion)Auto Loader(自動ローダー) の使用を推奨しています。Auto Loader はクラウドオブジェクトストレージに到着した新しいファイルを自動的に検出・処理します。

  • format("cloudFiles") を指定し、option("cloudFiles.format", "json") などで元データ形式を指定する。
  • option("cloudFiles.schemaLocation", ...) でスキーマ情報の保存先を指定する。
  • 構造化ストリーミング(readStream / writeStream)を用い、option("checkpointLocation", ...) で処理位置を記録する。
  • trigger(availableNow=True)(Scala では Trigger.AvailableNow)で「今ある新規ファイルをまとめて処理して停止」する。
  • toTable(table_name) で Delta テーブルへ書き込む。

Databricks は保存先に Delta Lake を推奨。Delta Lake は ACID トランザクションを提供し、レイクハウスを実現するオープンソースのストレージレイヤー。作成されるテーブルの既定形式が Delta Lake。

変換(Transform)と読み込み(Load)

  • 読み込み: spark.read.table(table_name) でテーブルを DataFrame として取得。
  • プレビュー: display(df) で対話的に表示。
  • 書き込み: バッチなら df.write...、ストリームなら writeStream...toTable()

変換(transformation)とアクション(action)・遅延評価(lazy evaluation)

Spark の DataFrame 操作は 2 種類に分かれます。

  • 変換(transformation): select, filter, withColumn, join, groupBy など、DataFrame から新しい DataFrame を作る操作。遅延評価され、この時点では計算は走らない(実行計画が構築されるだけ)。
  • アクション(action): display, count, collect, write / writeStream / toTable など、実際に計算を起動して結果を返す/書き出す操作。アクションが呼ばれて初めて、それまでの変換がまとめて最適化・実行される。

この「変換は貯めておき、アクションでまとめて実行する」モデルが遅延評価であり、Spark がクエリ全体を見て最適化できる理由です。

代替: Lakeflow パイプライン(旧 Delta Live Tables)を使うと、運用 ETL パイプラインの構築・デプロイ・保守の複雑さを軽減できる。手書きの構造化ストリーミングコードの代わりに宣言的にパイプラインを定義できる。


3-2. PySpark DataFrame API の主要操作(select / filter / withColumn / join / groupBy 等)

DataFrame は名前付き列を持つ分散データセットで、以下の変換をメソッドチェーンでつないで記述します(いずれも遅延評価される変換)。

操作意味
select列の選択・射影df.select("*", col("_metadata.file_path").alias("source_file"))
filter / where行のフィルタリングdf.filter(col("score") > 20)
withColumn列の追加・置換df.withColumn("name_length", get_name_length(df.name))
alias列やテーブルの別名付けcol("...").alias("source_file")
joinテーブル同士の結合df1.join(df2, "key")
groupBy + aggグループ化と集計df.groupBy("name_length").agg(total_score_udf(df["score"]).alias("total_score"))
display結果の対話的表示(アクション)display(df)

補足:

  • 列参照col("列名")from pyspark.sql.functions import col)や df["列名"]df.列名、Scala では $"列名" などで行う。
  • 組み込み関数(例: current_timestamp())は pyspark.sql.functions からインポートして使う。
  • 新しい列を計算列として付けるときは withColumn("新列名", 式)
  • select("*", ...) のように *(すべての列)に追加列を並べて射影できる。

ETL チュートリアルの変換例(メタデータ列と処理時刻の付与):

python
.select("*", col("_metadata.file_path").alias("source_file"), current_timestamp().alias("processing_time"))

UDF ページの集計例(グループ化して集計):

python
result_df = (df.groupBy("name_length")
  .agg(total_score_udf(df["score"]).alias("total_score")))
display(result_df)

重要な設計原則(メダリオン由来): 大規模データや定期・継続実行のワークロード(ETL・ストリーミング)では、組み込みの Apache Spark 関数(分散処理向けに最適化されている)を優先する。UDF はアドホッククエリ・手動クレンジング・探索的分析・小〜中規模データに向く。


3-3. Spark SQL / SQL 言語(テーブル・ビュー・CTAS)

Databricks SQL 言語リファレンスは、SQL コマンドを次のカテゴリに整理しています。

  • 一般リファレンス: データ型、関数、識別子(Identifiers)、リテラル、NULL のセマンティクス、パーティション、照合順序(Collation)など。
  • DDL ステートメント(Data Definition Language): オブジェクトの構造を作成・変更。ALTER, CREATE, DROP 系。
  • DML ステートメント(Data Manipulation Language): Delta Lake テーブルのデータを追加・変更・削除。COPY INTO, DELETE FROM, INSERT, LOAD DATA, MERGE INTO, UPDATE
  • データ取得ステートメント: SELECT(サブセレクト)、VALUES 句、SQL パイプライン構文、EXPLAIN。標準の SELECT 構文と SQL パイプライン構文の両方をサポート。
  • クエリ句: SELECT, WHERE, GROUP BY, HAVING, QUALIFY, ORDER BY, LIMIT, OFFSET, JOIN, PIVOT / UNPIVOT, CTE(共通テーブル式), 集合演算子(UNION / INTERSECT / EXCEPT)など。
  • Delta Lake ステートメント: OPTIMIZE, VACUUM, DESCRIBE HISTORY, RESTORE, CONVERT TO DELTA, CACHE SELECT など。
  • 補助ステートメント: DESCRIBE ..., SHOW ..., USE CATALOG / USE SCHEMA, SET, ANALYZE TABLE など。
  • セキュリティステートメント: GRANT, REVOKE, DENY, SHOW GRANTS など。

CREATE TABLE — マネージドテーブルと外部テーブル、CTAS

テーブルは既存スキーマ内に定義します。目的に応じて複数の作成方法があります。

  • CREATE TABLE [USING]: 最も一般的。次の 3 パターンで使う。
    • 列定義を明示して作る。
    • 既存ストレージ場所のデータから派生させる(→ 外部テーブルになり得る)。
    • クエリから派生させる = CTAS(CREATE TABLE AS SELECT)
  • CREATE TABLE (Hive 形式): Hive 構文版。CREATE TABLE [USING] の使用が推奨される。
  • CREATE TABLE LIKE: 別テーブルのデータではなく定義(スキーマ)だけをコピーして新テーブルを作る。
  • CREATE TABLE CLONE: Delta テーブルの複製。
    • DEEP CLONE(ディープクローン): 定義とデータを含む完全かつ独立したコピーを作成。
    • シャロークローン(浅いクローン): 定義のみコピーし、初期データは元テーブルのストレージを参照。どちらか一方を更新しても他方に影響しない(ただし新テーブルはソースの存在と列定義に依存)。

マネージドテーブル(managed table)と外部テーブル(external table / unmanaged table)の違い(試験頻出):

観点マネージドテーブル外部テーブル
データの所有・管理Databricks がデータとメタデータの両方を管理メタデータは Databricks、データ本体は外部の既存ストレージ場所(LOCATION 指定)
DROP 時の挙動メタデータとデータ本体の両方が削除されるメタデータのみ削除、データ本体は残る
作り方の目安列定義やクエリから作成(場所指定なし)既存のストレージ場所を指す LOCATION を指定して作成
MERGE のターゲット不可MERGE INTO のターゲットに外部テーブルは指定できない)

テーブルの既定フォーマットは Delta Lake

CREATE VIEW — ビュー・一時ビュー

ビュー(view)は SQL クエリの結果セットに基づく、物理データを持たない仮想テーブルです。ALTER VIEW / DROP VIEW はメタデータのみを変更します。

  • CREATE [OR REPLACE] VIEW: OR REPLACE は同名ビューを置換(DROP VIEW IF EXISTS + CREATE VIEW と等価)。ただし置換すると元の権限は保持されない。
  • TEMPORARY VIEW(一時ビュー): 作成したセッション内でのみ可視で、セッション終了時に削除される。名前を修飾(カタログ.スキーマ.名前)できない。
  • GLOBAL TEMPORARY VIEW(グローバル一時ビュー): システムが保持する一時スキーマ global_temp に関連付けられる(Databricks Runtime)。
  • IF NOT EXISTS: 既に存在すれば作成をスキップ。IF NOT EXISTSOR REPLACE は同時指定不可。
  • メトリックビュー(metric view): LANGUAGE YAMLWITH METRICS で定義する、ディメンションとメジャーを持つビュー(パブリックプレビュー、Unity Catalog 限定)。
  • スキーマバインディング(schema_binding): 基礎オブジェクトの変更にビューがどう追従するかを指定。SCHEMA BINDING(変更で無効化。ただし追加列は無視、安全な縮小キャストは許容)/ SCHEMA COMPENSATION(既定。ANSI キャスト可能なら許容)/ SCHEMA [TYPE] EVOLUTION(型変更・列追加削除を取り込む)。

ビューとテーブルの違い: テーブルは物理データを持つが、ビューは物理データを持たずクエリ定義を保持するだけ。ビューは複雑なクエリの再利用・抽象化・アクセス制御に使う。


3-4. UDF(ユーザー定義関数)の種類と使いどころ・性能特性

UDF(user-defined function) は、Databricks の組み込み機能を拡張し、複雑な計算・変換・カスタムデータ操作を再利用・共有するための関数です。

UDF を使うべき場面 / 使うべきでない場面

  • UDF が向く: 組み込みの Apache Spark 関数では表現しにくいロジック。アドホッククエリ、手動データクレンジング、探索的データ分析、小〜中規模データセット。代表例はデータの暗号化・復号・ハッシュ、JSON 解析、検証。
  • 組み込み関数 / Spark メソッドが向く: 大規模データセット、ETL ジョブやストリーミングなど定期・継続実行のワークロード。組み込み Apache Spark 関数は分散処理向けに最適化され、大規模で高性能。

UDF の種類

種類English入出力特徴
スカラー UDFscalar UDF1 行 → 1 値各行に 1 つの結果値。Unity Catalog 管理 or セッションスコープ。
バッチスカラー UDFbatch scalar UDFバッチ → 1:1 の行数行をバッチ処理し行ごとオーバーヘッドを削減。バッチ間で状態を保持。
非スカラー UDFnon-scalar UDF1:N / 多:多列/データセット全体で動作。pandas UDF(Series→Series ほか)が含まれる。
UDAFuser-defined aggregate function複数行 → 1 集計値集計を返す。セッションスコープ限定。
UDTFuser-defined table function1 引数以上 → 複数行(複数列可)入力行ごとに複数行を返す。Unity Catalog 管理 or セッションスコープ。

Unity Catalog 管理 UDF とセッションスコープ UDF

考慮事項Unity Catalog 管理 UDFセッションスコープ UDF
最適な用途チーム・ノートブック・ジョブ・SQL ウェアハウス間での安全な共有1 ノートブック/ジョブ内での迅速で反復的な開発
対応言語SQL, Python, Scala, JavaSQL, Python, Scala
ガバナンス・共有Unity Catalog の権限で管理、カタログエクスプローラーで検出可能現在の SparkSession に限定。管理・共有されない
永続性Unity Catalog に永続化、セッション間で再利用可現在のセッションにのみ存在

性能特性(試験頻出)

  • 組み込み関数と 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 は Apache Arrow を使いシリアライズコストを削減するため、行ごとの Python UDF より最大 100 倍高速。ベクトル化操作をサポート。

性能の目安(速い順): 組み込み関数 / SQL UDF > Scala UDF > pandas UDF(Arrow でベクトル化)> 行ごとの Python UDF


3-5. MERGE INTO によるアップサート(WHEN MATCHED / WHEN NOT MATCHED / WHEN NOT MATCHED BY SOURCE)

MERGE INTO は、ソーステーブル(source)に基づいて、ターゲット Delta テーブル(target)へ更新(UPDATE)・挿入(INSERT)・削除(DELETE)をまとめて適用する文です。これにより アップサート(upsert) を 1 文で実現できます。

重要な制約: MERGEDelta Lake テーブルでのみサポート。ターゲットは Delta テーブルでなければならず、外部テーブル(external table)はターゲットにできない。ターゲット名に options 仕様を含めてはならない。

構文

[ common_table_expression ]
  MERGE [ WITH SCHEMA EVOLUTION ] INTO target_table_name [target_alias]
     USING source_table_reference [source_alias]
     ON merge_condition
     { WHEN MATCHED [ AND matched_condition ] THEN matched_action |
       WHEN NOT MATCHED [BY TARGET] [ AND not_matched_condition ] THEN not_matched_action |
       WHEN NOT MATCHED BY SOURCE [ AND not_matched_by_source_condition ] THEN not_matched_by_source_action } [...]

matched_action
 { DELETE |
   UPDATE SET * [ EXCEPT ( column [, ...] ) ] |
   UPDATE SET { column = { expr | DEFAULT } } [, ...] }

not_matched_action
 { INSERT * [ EXCEPT ( column [, ...] ) ] |
   INSERT (column1 [, ...] ) VALUES ( expr | DEFAULT ] [, ...] )

not_matched_by_source_action
 { DELETE |
   UPDATE SET { column = { expr | DEFAULT } } [, ...] }

3 つの WHEN 句

  • WHEN MATCHED: ON merge_condition(+任意の matched_condition)でソース行がターゲット行と一致する場合に実行。
    • アクション: DELETE(一致するターゲット行を削除)/ UPDATE SET *(ソースの対応列でターゲット全列を更新。列名が同じであることが前提。異なると分析エラー)/ UPDATE SET col = expr(列ごとに更新)。
    • UPDATE SET * EXCEPT (col, ...): 特定列を除外して更新。除外列は null に設定される(DBR 12.2 LTS 以上)。
    • DEFAULT を式に指定して列を既定値に更新できる(DBR 11.3 LTS 以上)。
    • WHEN MATCHED 句を複数書ける。指定した順に評価され、最後の句以外は必ず matched_condition を持つ必要がある(さもないと NON_LAST_MATCHED_CLAUSE_OMIT_CONDITION エラー)。
  • WHEN NOT MATCHED [BY TARGET]: ソース行がどのターゲット行にも一致しない場合に行を挿入WHEN NOT MATCHED BY TARGETWHEN NOT MATCHED の別名、DBR 12.2 LTS 以上)。
    • アクション: INSERT *(ソースの対応列でターゲット全列を挿入。列が同じ前提)/ INSERT (col, ...) VALUES (expr, ...)(指定列のみ。未指定列は既定値または NULL)。
    • INSERT * EXCEPT (col, ...): 特定列を除外(除外列は null)。
    • 複数句は順に評価。最後以外は not_matched_condition が必須
  • WHEN NOT MATCHED BY SOURCE(DBR 12.2 LTS 以上): ターゲット行がソースのどの行にも一致しない場合に実行。ソースから消えたレコードをターゲットで削除/更新する用途。
    • アクション: DELETE / UPDATE SET col = expr。式はターゲットテーブルの列のみ参照可(not_matched_by_source_condition もターゲット列のみ参照するブール式)。
    • 注意: 条件なしだと大量のターゲット行が変更され得る。パフォーマンスのため not_matched_by_source_condition で対象を限定すること。複数句は順に評価、最後以外は条件必須。

重要な注意点

  • 複数一致(multiple matches): ONWHEN MATCHED の条件により、ソースの複数行がターゲットの同一行に一致すると操作は失敗しエラーになる(どのソース行で更新すべきか曖昧なため)。対策としてソースを前処理し、各キーの最新変更のみを残す(CDC の例参照)。ただし無条件の DELETE は複数一致でも曖昧ではないため許可される。
  • WITH SCHEMA EVOLUTION(DBR 15.2 以上): 自動スキーマ進化を有効化。ターゲット Delta テーブルのスキーマがソースに合わせて自動更新される。
  • 共通テーブル式(CTE) を先頭に付けられる。

代表的な用途

MERGE INTO は次のような複雑な操作に使えます。

  • データの重複除去(deduplication)
  • 変更データのアップサート(CDC の適用)
  • SCD Type 2(Slowly Changing Dimension Type 2)操作の適用(履歴を新しい行として保持)

3-6. メダリオンアーキテクチャ(Bronze / Silver / Gold の役割と設計指針)

メダリオンアーキテクチャ(medallion architecture / マルチホップアーキテクチャ multi-hop architecture) は、レイクハウス内のデータを論理的に整理し、層を経るごとに構造と品質を段階的に向上させるデータ設計パターンです。Bronze(未加工)→ Silver(検証済み)→ Gold(エンリッチ済み)と進めることで、ACID を保証しつつ、BI や ML に適した信頼できる単一ソースを構築します。推奨されるが必須ではありません

質問Bronze(銅)Silver(銀)Gold(金)
何が起こるか生データ取り込みデータのクリーニングと検証ディメンションモデリングと集計
主な利用者データエンジニア、データオペレーション、コンプライアンス/監査チームデータエンジニア、データアナリスト、データサイエンティストビジネスアナリスト/BI 開発者、データサイエンティスト/ML エンジニア、役員・意思決定者、運用チーム

例では、各層を別スキーマに格納する(例: ops.bronze, ops.silver, ops.gold)。

Bronze 層(生データの取り込み)

  • データソースの生の状態を元の形式で保存・維持する。クレンジングや検証はほぼ行わない(最小限の検証のみ)。
  • 増分的に追加され、時間とともに増加する。すべての履歴データを保持し、再処理と監査を可能にする。
  • 信頼の単一ソース(single source of truth) として機能。データ忠実性を維持。
  • アナリストやデータサイエンティストが直接アクセスするのではなく、Silver を強化するワークロードが使うことを想定。
  • ソースはクラウドオブジェクトストレージ(S3, GCS, ADLS)、メッセージバス(Kafka, Kinesis)、フェデレーションシステム(Lakehouse フェデレーション)など。ストリーミングとバッチを任意に組み合わせ可能。
  • 設計のコツ: データ欠落を防ぎ予期しないスキーマ変更から保護するため、ほとんどのフィールドを STRING、VARIANT、または BINARY として格納することを推奨。データの実績やソースなどのメタデータ列(例: _metadata.file_name)を追加することがある。

Silver 層(検証・重複排除)

  • 1 つ以上の Bronze/Silver テーブルから読み取り、Silver テーブルへ書き込むインジェストから直接 Silver に書くのは非推奨(スキーマ変更やレコード破損でエラーになるため)。
  • すべてのソースが追加専用(append-only)と仮定し、Bronze からの読み取りはほとんどストリーミング読み取りで構成。バッチ読み取りは小さなデータセット(小さなディメンションテーブル)に限定。
  • 各レコードの、少なくとも 1 つの検証済み・非集計表現を必ず含む。集計表現は通常 Gold に置く(多数のダウンストリームを駆動する場合は Silver に置くこともある)。
  • ここで実施する処理: スキーマ適用、null/欠損値の処理、重複排除(deduplication)、順序が乱れた/遅れて到着するデータの解決、データ品質チェックと適用、スキーマ展開、型キャスト、結合(join)、正規化
  • データモデリングを開始する層。ネスト/半構造化データの表現方法(VARIANT 型、JSON 文字列、struct/map/array、スキーマのフラット化、複数テーブルへの正規化)を選ぶ。

Gold 層(分析強化)

  • ダウンストリームの分析・ダッシュボード・ML・アプリを駆動する、高度に洗練されたビュー。特定の期間や地域に対して高度に集計・フィルタされることが多い。
  • 含まれるデータセット数は Silver / Bronze より少ない。ビジネス機能・ニーズに対応する、意味的に意味のあるデータセット。
  • ビジネスロジックと要件に合わせる。ディメンションモデル(dimensional model)でリレーションシップを確立しメジャーを定義。人事・財務・IT など複数の Gold 層を作ることもある。
  • 集計を作成(平均・カウント・最大・最小など)。頻繁に使う集計はマテリアライズドビュー(materialized view) として事前集計しておく。
  • クエリ/ダッシュボード性能を最適化(頻繁に照会されるため)。大量の履歴データは通常 Silver でアクセスし、Gold には具体化しない。

Gold 層の集計例(週次売上のマテリアライズドビュー):

sql
CREATE OR REPLACE MATERIALIZED VIEW weekly_sales AS
SELECT week,
       prod_id,
       region,
       SUM(units) AS total_units,
       SUM(units * rate) AS total_sales
FROM orders
GROUP BY week, prod_id, region

データインジェスト頻度とコストのトレードオフ

インジェスト頻度コスト遅延(レイテンシ)
継続的な増分インジェスト高い低いspark.readStream で継続実行するストリーミングテーブル/構造化ストリーミング
トリガーされた増分インジェスト低い高いspark.readStream + Trigger.AvailableNow、スケジュール/ファイル到着トリガー
手動増分を伴うバッチインジェスト低い最も高い(実行頻度が低い)spark.read、パーティション上書きなどのプリミティブ、datetime でのパーティション分割

4. 構文・コード例

4-1. Auto Loader で Delta テーブルへ増分インジェスト(PySpark)

python
# Import functions
from pyspark.sql.functions import col, current_timestamp

# Define variables used in code below
file_path = "/databricks-datasets/structured-streaming/events"
username = spark.sql("SELECT regexp_replace(session_user(), '[^a-zA-Z0-9]', '_')").first()[0]
table_name = f"{username}_etl_quickstart"
checkpoint_path = f"/tmp/{username}/_checkpoint/etl_quickstart"

# Clear out data from previous demo execution
spark.sql(f"DROP TABLE IF EXISTS {table_name}")
dbutils.fs.rm(checkpoint_path, True)

# Configure Auto Loader to ingest JSON data to a Delta table
(spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "json")
  .option("cloudFiles.schemaLocation", checkpoint_path)
  .load(file_path)
  .select("*", col("_metadata.file_path").alias("source_file"), current_timestamp().alias("processing_time"))
  .writeStream
  .option("checkpointLocation", checkpoint_path)
  .trigger(availableNow=True)
  .toTable(table_name))

4-2. DataFrame の読み込みとプレビュー

python
df = spark.read.table(table_name)
display(df)

4-3. スカラー UDF(SQL 版と PySpark 版)

SQL UDF(Unity Catalog に登録):

sql
-- Create a SQL UDF for name length
CREATE OR REPLACE FUNCTION main.test.get_name_length(name STRING)
RETURNS INT
RETURN LENGTH(name);

-- Use the UDF in a SQL query
SELECT name, main.test.get_name_length(name) AS name_length
FROM your_table;

PySpark スカラー UDF:

python
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))

# Show the result
display(df)

4-4. pandas UDF(Series → Series、ベクトル化)

python
from pyspark.sql.functions import pandas_udf
import pandas as pd

df = spark.createDataFrame([(70, 1.75), (80, 1.80), (60, 1.65)], ["Weight", "Height"])

@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()

4-5. UDAF(pandas UDF による集計)

python
from pyspark.sql.functions import pandas_udf
from pyspark.sql import SparkSession
import pandas as pd

# Define a pandas UDF for aggregating scores
@pandas_udf("int")
def total_score_udf(scores: pd.Series) -> int:
  return scores.sum()

# Group by name length and aggregate
result_df = (df.groupBy("name_length")
  .agg(total_score_udf(df["score"]).alias("total_score")))

display(result_df)

4-6. UDTF(入力行ごとに複数行を返す)

SQL 版:

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);

PySpark 版:

python
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()

4-7. MERGE INTO(アップサートの各パターン)

WHEN MATCHED(削除・条件付き更新・複数句):

sql
-- Delete all target rows that have a match in the source table.
MERGE INTO target USING source
  ON target.key = source.key
  WHEN MATCHED THEN DELETE

-- Conditionally update target rows that have a match in the source table using the source value.
MERGE INTO target USING source
  ON target.key = source.key
  WHEN MATCHED AND target.updated_at < source.updated_at THEN UPDATE SET *

-- Multiple MATCHED clauses: conditionally delete, otherwise update two columns.
MERGE INTO target USING source
  ON target.key = source.key
  WHEN MATCHED AND target.marked_for_deletion THEN DELETE
  WHEN MATCHED THEN UPDATE SET target.updated_at = source.updated_at, target.value = DEFAULT

WHEN NOT MATCHED [BY TARGET](挿入):

sql
-- Insert all rows from the source that are not already in the target table.
MERGE INTO target USING source
  ON target.key = source.key
  WHEN NOT MATCHED THEN INSERT *

-- Conditionally insert new rows using unmatched source rows.
MERGE INTO target USING source
  ON target.key = source.key
  WHEN NOT MATCHED BY TARGET AND source.created_at > now() - INTERVAL '1' DAY
    THEN INSERT (created_at, value) VALUES (source.created_at, DEFAULT)

-- Insert all columns except `last_updated`, which is set to null in the target.
MERGE INTO target USING source
  ON target.key = source.key
  WHEN NOT MATCHED THEN INSERT * EXCEPT (last_updated)

WHEN NOT MATCHED BY SOURCE(ソースにない行の削除・更新):

sql
-- Delete all target rows that have no matches in the source table.
MERGE INTO target USING source
  ON target.key = source.key
  WHEN NOT MATCHED BY SOURCE THEN DELETE

-- Multiple NOT MATCHED BY SOURCE clauses.
MERGE INTO target USING source
  ON target.key = source.key
  WHEN NOT MATCHED BY SOURCE AND target.marked_for_deletion THEN DELETE
  WHEN NOT MATCHED BY SOURCE THEN UPDATE SET target.value = DEFAULT

WITH SCHEMA EVOLUTION(3 句をまとめて、スキーマ進化あり):

sql
MERGE WITH SCHEMA EVOLUTION INTO target USING source
  ON source.key = target.key
  WHEN MATCHED THEN UPDATE SET *
  WHEN NOT MATCHED THEN INSERT *
  WHEN NOT MATCHED BY SOURCE THEN DELETE

4-8. CREATE VIEW(通常ビュー・一時ビュー)

sql
-- Create or replace view with column comments.
CREATE OR REPLACE VIEW experienced_employee
    (id COMMENT 'Unique identification number', Name)
    COMMENT 'View for experienced employees'
    AS SELECT id, name
         FROM all_employee
        WHERE working_years > 5;

-- Create a temporary view.
CREATE TEMPORARY VIEW subscribed_movies
    AS SELECT mo.member_id, mb.full_name, mo.movie_title
         FROM movies AS mo
         INNER JOIN members AS mb
            ON mo.member_id = mb.id;

5. 試験で問われるポイント

  • 変換とアクションの区別、遅延評価: select/filter/withColumn/join/groupBy は変換で遅延評価。count/collect/display/write/toTable はアクションで実行が起動する。「いつ実際に計算が走るか」を理解する。
  • Auto Loader: 増分インジェストの推奨手段。format("cloudFiles")cloudFiles.formatcloudFiles.schemaLocationcheckpointLocationtrigger(availableNow=True) の役割。
  • 既定フォーマットは Delta Lake、ACID を提供。MERGE は Delta テーブル限定・外部テーブルはターゲット不可。
  • マネージドテーブル vs 外部テーブル: DROP 時にデータが消えるか残るか。外部テーブルは LOCATION を持ち、データ本体は外部管理。
  • CTAS(CREATE TABLE AS SELECT)CREATE TABLE LIKE(定義のみ)、CLONE(DEEP/浅い)の違い。
  • ビューの種類: 通常ビュー(物理データなし)、TEMPORARY VIEW(セッション限定)、GLOBAL TEMPORARY VIEWglobal_temp)。CREATE OR REPLACE VIEW は権限を引き継がない。
  • UDF の性能順: 組み込み関数 / SQL UDF > Scala UDF > pandas UDF > 行ごと Python UDF。pandas UDF は Arrow でベクトル化し Python UDF の最大 100 倍。大規模/本番は組み込み関数を優先。
  • UDF の分類: スカラー / バッチスカラー / 非スカラー(pandas)/ UDAF / UDTF。Unity Catalog 管理(共有・永続・ガバナンス)とセッションスコープ(一時・現在の SparkSession 限定)の違い。
  • MERGE INTO の 3 句: WHEN MATCHED(UPDATE/DELETE)、WHEN NOT MATCHED [BY TARGET](INSERT)、WHEN NOT MATCHED BY SOURCE(UPDATE/DELETE)。UPDATE SET * / INSERT * は列名一致が前提。複数句は順に評価し最後以外は条件必須。複数一致はエラー(前処理で回避)。
  • MERGE のユースケース: アップサート、重複排除、CDC 適用、SCD Type 2。
  • メダリオンアーキテクチャ: Bronze(生・未検証・STRING/VARIANT/BINARY 推奨・履歴保持)→ Silver(洗浄・検証・重複排除・結合・モデリング開始・Bronze からはストリーミング読み取り推奨)→ Gold(集計・ディメンションモデル・マテリアライズドビュー・少数の洗練データセット)。
  • インジェスト頻度のトレードオフ: 継続(高コスト・低遅延)/トリガー(低コスト・高遅延)/バッチ(低コスト・最高遅延)。

6. 理解度チェックリスト

  • [ ] 変換(transformation)とアクション(action)を具体例で区別でき、遅延評価がいつ実行を起動するか説明できる
  • [ ] Auto Loader(cloudFiles)で JSON を Delta テーブルへ増分インジェストするコードを書ける
  • [ ] checkpointLocationtrigger(availableNow=True) の役割を説明できる
  • [ ] Databricks の既定テーブル形式が Delta Lake であり、ACID を提供することを説明できる
  • [ ] select / filter / withColumn / join / groupBy + agg の役割と書き方を説明できる
  • [ ] マネージドテーブルと外部テーブルの違い(特に DROP 時の挙動と MERGE ターゲット可否)を説明できる
  • [ ] CTAS、CREATE TABLE LIKECLONE(DEEP / 浅い)の違いを説明できる
  • [ ] 通常ビュー・一時ビュー・グローバル一時ビューの違いと有効範囲を説明できる
  • [ ] CREATE OR REPLACE VIEW が権限を保持しないことを知っている
  • [ ] UDF を使うべき場面と、組み込み関数を優先すべき場面を判断できる
  • [ ] スカラー / バッチスカラー / 非スカラー(pandas) / UDAF / UDTF を区別できる
  • [ ] Unity Catalog 管理 UDF とセッションスコープ UDF の違いを説明できる
  • [ ] UDF の性能順(組み込み/SQL > Scala > pandas > 行ごと Python)と pandas UDF が Arrow で高速な理由を説明できる
  • [ ] MERGE INTO の基本構文と、WHEN MATCHED / WHEN NOT MATCHED [BY TARGET] / WHEN NOT MATCHED BY SOURCE の各アクションを書ける
  • [ ] UPDATE SET * / INSERT * の前提(列名一致)と EXCEPT の挙動を説明できる
  • [ ] MERGE の複数一致がエラーになる理由と、CDC 前処理での回避策を説明できる
  • [ ] MERGE でアップサート・重複排除・SCD Type 2 を実現できることを説明できる
  • [ ] メダリオンアーキテクチャの Bronze / Silver / Gold の役割・利用者・データ状態を説明できる
  • [ ] Bronze で STRING/VARIANT/BINARY 保存やメタデータ列付与を推奨する理由を説明できる
  • [ ] Silver で Bronze からストリーミング読み取りを推奨する理由を説明できる
  • [ ] Gold でマテリアライズドビューやディメンションモデルを使う目的を説明できる
  • [ ] インジェスト頻度(継続 / トリガー / バッチ)のコストと遅延のトレードオフを説明できる