Skip to content

Free Edition ハンズオン課題|Associate 対策(全12課題)

Databricks 認定データエンジニア Associate の 7 セクションに対応した、Free Edition で実際に動かせる課題集です。 先に 00-free-edition.md のセットアップ(スキーマ de_sandbox とボリューム landing の作成)を済ませてください。 各課題は「目的 → 対応教材 → 手順・コード → 確認ポイント」の構成です。確認ポイントは試験で問われる形にしてあります。

前提: 以下のコードはすべて <catalog> を自分のワークスペースカタログ名に置き換えて実行します。ノートブックの先頭で一度実行しておくと以降が短く書けます。

sql
USE CATALOG <catalog>;
USE SCHEMA de_sandbox;
python
CATALOG = "<catalog>"
SCHEMA  = "de_sandbox"
VOL     = f"/Volumes/{CATALOG}/{SCHEMA}/landing"

課題一覧

#課題対応セクション所要
1Delta の ACID・タイムトラベル・トランザクションログを覗く20分
2マネージド/外部テーブルとスキーマ適用・進化① ③20分
3Auto Loader で増分取り込み(スキーマ推論・進化・_rescued_data40分
4COPY INTO と冪等性・VALIDATE20分
5PySpark DataFrame API と変換/アクション・遅延評価30分
6UDF 3 種の性能を実測する(組み込み / pandas / 行ごと Python)30分
7MERGE INTO でアップサート・重複排除・SCD Type 240分
8メダリオン(Bronze→Silver→Gold)を一気通貫で作る40分
9Lakeflow Jobs で 3 タスクの DAG・cron・通知・Repair run40分
10宣言型パイプラインと Expectations(warn / drop / fail)40分
11Unity Catalog の権限と USE CATALOG の壁・ボリューム30分
12OPTIMIZE / 液体クラスタリング / VACUUM とクエリプロファイル40分
Git フォルダーと Databricks CLI・バンドル40分

課題 1|Delta の ACID・タイムトラベル・トランザクションログを覗く

目的: 「書き込みごとにバージョンが増える」「DROP してもタイムトラベルで戻れる」を体で覚える。 対応教材: 01-platform.md 3-3 / 確認問題: 01-platform.md Q1〜Q4

sql
CREATE OR REPLACE TABLE t_orders (id INT, amount DECIMAL(10,2), status STRING);

INSERT INTO t_orders VALUES (1, 100.00, 'NEW');           -- v1
INSERT INTO t_orders VALUES (2, 250.00, 'NEW');           -- v2
UPDATE t_orders SET status = 'PAID' WHERE id = 1;         -- v3
DELETE FROM t_orders WHERE id = 2;                        -- v4

-- 1) 書き込みごとにバージョンが増えることを確認
DESCRIBE HISTORY t_orders;

-- 2) タイムトラベル:削除前に戻って読む
SELECT * FROM t_orders VERSION AS OF 2;
SELECT * FROM t_orders TIMESTAMP AS OF (SELECT timestamp FROM (DESCRIBE HISTORY t_orders) WHERE version = 2);

-- 3) 物理構成を見る(ファイル数・サイズ・場所)
DESCRIBE DETAIL t_orders;

-- 4) 復元
RESTORE TABLE t_orders TO VERSION AS OF 2;
SELECT * FROM t_orders;
DESCRIBE HISTORY t_orders;   -- RESTORE も新しいバージョンとして記録される

確認ポイント

  • DESCRIBE HISTORYoperation 列に WRITE / UPDATE / DELETE / RESTORE が並ぶ。コミットごとにバージョンがインクリメントされるのがトランザクションログの役割
  • RESTORE 自体も新バージョンになる(履歴は消えない)
  • 「原子性=完全に成功か完全に失敗」「永続性=コミット済みは永続」を自分の言葉で説明できるか

課題 2|マネージド/外部テーブルとスキーマ適用・進化

目的: DROP 時の挙動の違い、スキーマ適用(拒否)とスキーマ進化(受け入れ) の差を確認する。 対応教材: 01-platform.md 3-3 / 03-transformation.md 3-3

python
# スキーマ適用:定義外の列を含む書き込みは拒否される
spark.sql("CREATE OR REPLACE TABLE t_schema (id INT, name STRING)")
spark.sql("INSERT INTO t_schema VALUES (1, 'a')")

df_extra = spark.createDataFrame([(2, "b", "extra")], "id INT, name STRING, memo STRING")

try:
    df_extra.write.mode("append").saveAsTable(f"{CATALOG}.{SCHEMA}.t_schema")
except Exception as e:
    print("拒否された(スキーマ適用):", type(e).__name__)

# スキーマ進化:mergeSchema で列を追加して受け入れる
(df_extra.write.mode("append")
    .option("mergeSchema", "true")
    .saveAsTable(f"{CATALOG}.{SCHEMA}.t_schema"))

spark.table(f"{CATALOG}.{SCHEMA}.t_schema").printSchema()
display(spark.table(f"{CATALOG}.{SCHEMA}.t_schema"))
sql
-- 外部テーブル相当の挙動はボリューム上のファイルで確認する
-- (Free Edition では外部ロケーション/ストレージ資格情報を作れないため、
--   「マネージドは DROP でデータも消える」ことだけ実機で確認する)
CREATE OR REPLACE TABLE t_managed AS SELECT * FROM samples.nyctaxi.trips LIMIT 100;
DESCRIBE DETAIL t_managed;   -- location がマネージドストレージ配下であることを確認
DROP TABLE t_managed;
-- 保持期間後にデータファイルも削除される(マネージドテーブルの挙動)

確認ポイント

  • スキーマ適用は書き込み時にスキーマを検証して要件に合わないデータを弾く。スキーマ進化はデータを書き換えずにスキーマを変更する。この 2 つを取り違えないこと
  • マネージド = DROP でデータも削除 / 外部 = メタデータのみ削除でデータは残る。さらに外部テーブルは MERGE INTO のターゲットにできない

課題 3|Auto Loader で増分取り込み(スキーマ推論・進化・_rescued_data

目的: cloudFiles の増分検知・チェックポイント・スキーマ進化・レスキュー列を一通り体験する。Associate 頻出対応教材: 02-ingestion.md確認問題: 02-ingestion.md

3-1. 取り込み元に JSON を 2 ファイル置く

python
import json
src  = f"{VOL}/events"
ckpt = f"{VOL}/_ckpt/events"

dbutils.fs.mkdirs(src)

batch1 = [{"id": 1, "type": "click", "ts": "2026-07-01T10:00:00"},
          {"id": 2, "type": "view",  "ts": "2026-07-01T10:00:05"}]
dbutils.fs.put(f"{src}/batch1.json", "\n".join(json.dumps(r) for r in batch1), overwrite=True)

3-2. Auto Loader で読み、AvailableNow で 1 回だけ処理して停止

python
from pyspark.sql.functions import col, current_timestamp

def ingest():
    q = (spark.readStream.format("cloudFiles")
        .option("cloudFiles.format", "json")
        .option("cloudFiles.schemaLocation", ckpt)        # スキーマ推論/進化の追跡場所
        .load(src)
        .select("*",
                col("_metadata.file_path").alias("source_file"),
                current_timestamp().alias("processing_time"))
        .writeStream
        .option("checkpointLocation", ckpt)               # 処理済みファイルの追跡=exactly-once の基盤
        .trigger(availableNow=True)                       # 今ある分を処理して停止(Free Edition はこれのみ)
        .toTable(f"{CATALOG}.{SCHEMA}.bronze_events"))
    q.awaitTermination()

ingest()
display(spark.table(f"{CATALOG}.{SCHEMA}.bronze_events"))

3-3. 2 回目を実行して「増分だけ」処理されることを確認

python
batch2 = [{"id": 3, "type": "click", "ts": "2026-07-01T10:01:00"}]
dbutils.fs.put(f"{src}/batch2.json", "\n".join(json.dumps(r) for r in batch2), overwrite=True)

ingest()   # batch1 は再処理されない
spark.sql(f"SELECT source_file, count(*) FROM {CATALOG}.{SCHEMA}.bronze_events GROUP BY 1").show(truncate=False)

3-4. 新しい列を追加してスキーマ進化を体験する

python
batch3 = [{"id": 4, "type": "purchase", "ts": "2026-07-01T10:02:00", "amount": 1200}]
dbutils.fs.put(f"{src}/batch3.json", "\n".join(json.dumps(r) for r in batch3), overwrite=True)

try:
    ingest()   # 既定 addNewColumns → UnknownFieldException で停止する
except Exception as e:
    print("スキーマ進化で停止:", str(e)[:200])

ingest()   # 再起動すると更新後のスキーマで処理が進む
spark.table(f"{CATALOG}.{SCHEMA}.bronze_events").printSchema()

3-5. _rescued_data を確認する

python
# 既定は全列 string 推論。型が合わない・スキーマにない値は _rescued_data へ退避される
(spark.readStream.format("cloudFiles")
    .option("cloudFiles.format", "json")
    .option("cloudFiles.schemaLocation", f"{VOL}/_ckpt/rescue")
    .option("cloudFiles.inferColumnTypes", "true")
    .option("cloudFiles.schemaEvolutionMode", "rescue")   # 進化させず全部レスキュー列へ
    .load(src)
    .writeStream
    .option("checkpointLocation", f"{VOL}/_ckpt/rescue")
    .trigger(availableNow=True)
    .toTable(f"{CATALOG}.{SCHEMA}.bronze_rescue")).awaitTermination()

spark.sql(f"SELECT _rescued_data FROM {CATALOG}.{SCHEMA}.bronze_rescue WHERE _rescued_data IS NOT NULL").show(truncate=False)

確認ポイント

  • cloudFiles.schemaLocation(スキーマ進化用)と checkpointLocation(処理済み追跡用)は役割が違う。試験で混同を狙われる
  • 既定のスキーマ進化モードは addNewColumns で、新列検出時に UnknownFieldException で停止 → 再起動で続行。だから本番では Lakeflow ジョブで自動再起動させる
  • JSON / CSV / XML の既定推論型はすべて文字列。型推論には cloudFiles.inferColumnTypes=true
  • Auto Loader はファイルの処理順序を保証しない
  • Free Edition の制約: 検出モードは既定のディレクトリ一覧のみ想定。ファイル通知モードはクラウドのイベント通知設定が必要なため実質不可(座学で補う)

課題 4|COPY INTO と冪等性・VALIDATE

目的: 「既読ファイルはスキップ」「force で冪等性無効化」「VALIDATE は書き込まない」を確認する。 対応教材: 02-ingestion.md 3-4 / 確認問題: Q7・Q8・Q10

sql
CREATE TABLE IF NOT EXISTS t_copy (id STRING, type STRING, ts STRING, amount STRING);

-- 1) 書き込まずに検証だけする
COPY INTO t_copy
FROM '/Volumes/<catalog>/de_sandbox/landing/events'
FILEFORMAT = JSON
VALIDATE ALL;
SELECT count(*) FROM t_copy;     -- 0 のまま

-- 2) 実際にロード
COPY INTO t_copy
FROM '/Volumes/<catalog>/de_sandbox/landing/events'
FILEFORMAT = JSON;
SELECT count(*) FROM t_copy;

-- 3) もう一度同じコマンド → 冪等なので 0 行(既読ファイルはスキップ)
COPY INTO t_copy
FROM '/Volumes/<catalog>/de_sandbox/landing/events'
FILEFORMAT = JSON;

-- 4) force で冪等性を無効化 → 既読ファイルも再読み込みされて件数が倍になる
COPY INTO t_copy
FROM '/Volumes/<catalog>/de_sandbox/landing/events'
FILEFORMAT = JSON
COPY_OPTIONS ('force' = 'true');
SELECT count(*) FROM t_copy;

確認ポイント

  • COPY INTO再試行可能かつ冪等。既読ファイルは内容が後から変わってもスキップされる
  • COPY_OPTIONS('force'='true')(既定 false)で冪等性が無効化される
  • FILES(最大 1000 ファイル)と PATTERN(glob)は併用不可
  • FORMAT_OPTIONS(リーダー設定)と COPY_OPTIONS(COPY INTO の動作)の違い

課題 5|PySpark DataFrame API と変換/アクション・遅延評価

目的: どの操作がジョブを起動するかを目で確認する。Associate/Professional 共通の最重要概念。 対応教材: 03-transformation.md 3-1・3-2 / 確認問題: Q1

python
from pyspark.sql.functions import col, desc, sum as _sum, count

# --- ここまでは変換だけ:ジョブは起動しない(ノートブック下部にジョブが出ないことを確認)---
df = spark.table("samples.nyctaxi.trips")
sub = (df
   .filter((col("trip_distance") > 1) & (col("fare_amount") > 0))   # narrow
   .withColumn("pickup_date", col("tpep_pickup_datetime").cast("date"))  # narrow
   .groupBy("pickup_date")                                          # wide → シャッフル
   .agg(_sum("fare_amount").alias("fare"), count("*").alias("n"))
   .orderBy(desc("fare")))                                          # wide → シャッフル

print("ここまでアクションなし")

# --- 実行計画を見る(シャッフル=Exchange の位置を確認)---
sub.explain("formatted")

# --- アクション:ここで初めてジョブが起動する ---
sub.show(5)          # アクション①
print(sub.count())   # アクション②(本番パイプラインでは避ける)
sub.write.mode("overwrite").saveAsTable(f"{CATALOG}.{SCHEMA}.gold_daily_fare")  # アクション③(本来これだけにすべき)

# --- 不変性の確認 ---
d1 = spark.range(3)
d2 = d1.withColumn("x", col("id") * 2)
d1.printSchema()   # 元の d1 は変わっていない
d2.printSchema()

確認ポイント

  • explain("formatted") に現れる Exchange がシャッフル。groupBy / orderBy / join の位置と一致する
  • filter / withColumn / selectnarrowgroupBy / join / orderBy / distinctwide
  • 運用パイプラインには書き込みアクションのみを置くべき理由(count / display が最適化を妨げる)を説明できるか
  • DataFrame は不変。変換結果は変数に代入しないと反映されない

課題 6|UDF 3 種の性能を実測する

目的: 「組み込み関数 > pandas UDF > 行ごと Python UDF」を自分の環境の数字で確認する。 対応教材: 03-transformation.md 3-4 / 確認問題: Q5

python
import time
import pandas as pd
from pyspark.sql.functions import udf, pandas_udf, col, sqrt as _sqrt, sum as _sum

# 注: Free Edition(サーバーレス)では df.cache() / persist() が例外になるためキャッシュは使わない
df = spark.range(0, 5_000_000).withColumn("v", col("id") * 1.0)

@udf("double")
def py_sqrt(v):                            # 行ごと Python UDF(最も遅い)
    return float(v) ** 0.5

@pandas_udf("double")
def pd_sqrt(s: pd.Series) -> pd.Series:    # pandas UDF(Arrow でベクトル化)
    return s ** 0.5

def bench(name, column):
    t = time.time()
    # agg(sum) は各行の値を必要とするため、UDF の評価が確実に走る
    df.select(column.alias("r")).agg(_sum("r")).collect()
    print(f"{name:20s}: {time.time() - t:6.2f} s")

bench("builtin sqrt",    _sqrt("v"))
bench("pandas UDF",      pd_sqrt("v"))
bench("python UDF (row)", py_sqrt("v"))

注: count() で計測すると列の値が不要と判断されて UDF の評価が省略されることがあります。上のように 値を使う集計sum)や Delta テーブルへの書き込みで強制してください。

確認ポイント

  • 教材の性能順「組み込み関数 / SQL UDF > Scala UDF > pandas UDF > 行ごとの Python UDF」のうち、Scala 以外は実測できる
  • pandas UDF が速い理由は Apache Arrow による列指向転送でシリアライズコストを削減し、ベクトル化するから(行ごと Python UDF の最大 100 倍)
  • Free Edition の制約: Scala UDF は非サポートなので実測不可。順序だけ暗記する
  • おまけ: UDF の NULL 落とし穴(WHERE x IS NOT NULL AND udf(x) は短絡保証がない)も試してみる

課題 7|MERGE INTO でアップサート・重複排除・SCD Type 2

目的: 3 つの WHEN 句と、複数一致でエラーになることを実際に踏む。 対応教材: 03-transformation.md 3-5 / 確認問題: Q3・Q4・Q9・Q10

sql
-- 準備
CREATE OR REPLACE TABLE dim_customer (id INT, name STRING, city STRING, updated_at TIMESTAMP);
INSERT INTO dim_customer VALUES
  (1, 'Alice', 'Tokyo',  '2026-07-01 00:00:00'),
  (2, 'Bob',   'Osaka',  '2026-07-01 00:00:00');

CREATE OR REPLACE TEMP VIEW updates AS SELECT * FROM VALUES
  (1, 'Alice', 'Kyoto',  TIMESTAMP'2026-07-02 00:00:00'),   -- 更新
  (3, 'Carol', 'Nagoya', TIMESTAMP'2026-07-02 00:00:00')    -- 挿入
  AS t(id, name, city, updated_at);

-- 1) 基本のアップサート(条件付き更新)
MERGE INTO dim_customer AS tgt
USING updates AS src ON tgt.id = src.id
WHEN MATCHED AND tgt.updated_at < src.updated_at THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *;
SELECT * FROM dim_customer ORDER BY id;

-- 2) ソースにない行を削除(必ず条件で範囲を絞る)
MERGE INTO dim_customer AS tgt
USING updates AS src ON tgt.id = src.id
WHEN NOT MATCHED BY SOURCE AND tgt.updated_at < TIMESTAMP'2026-07-02 00:00:00' THEN DELETE;
SELECT * FROM dim_customer ORDER BY id;
sql
-- 3) 複数一致でエラーになることを確認する(試験頻出)
CREATE OR REPLACE TEMP VIEW dup_updates AS SELECT * FROM VALUES
  (1, 'Alice', 'Sapporo', TIMESTAMP'2026-07-03 00:00:00'),
  (1, 'Alice', 'Fukuoka', TIMESTAMP'2026-07-04 00:00:00')   -- 同じキーが 2 行
  AS t(id, name, city, updated_at);

-- ↓ これは失敗する(どのソース行で更新すべきか曖昧)
MERGE INTO dim_customer AS tgt
USING dup_updates AS src ON tgt.id = src.id
WHEN MATCHED THEN UPDATE SET *;

-- 4) 回避策:ソースを前処理してキーごとに最新のみ残す
CREATE OR REPLACE TEMP VIEW dedup_updates AS
SELECT * FROM (
  SELECT *, row_number() OVER (PARTITION BY id ORDER BY updated_at DESC) AS rn
  FROM dup_updates
) WHERE rn = 1;

MERGE INTO dim_customer AS tgt
USING dedup_updates AS src ON tgt.id = src.id
WHEN MATCHED THEN UPDATE SET tgt.city = src.city, tgt.updated_at = src.updated_at;
SELECT * FROM dim_customer ORDER BY id;
sql
-- 5) insert-only マージによる冪等な重複排除
CREATE OR REPLACE TABLE logs (uniqueId STRING, msg STRING, date DATE);
CREATE OR REPLACE TEMP VIEW newLogs AS SELECT * FROM VALUES
  ('a', 'hello', DATE'2026-07-01'), ('b', 'world', DATE'2026-07-01') AS t(uniqueId, msg, date);

MERGE INTO logs USING newLogs ON logs.uniqueId = newLogs.uniqueId
WHEN NOT MATCHED THEN INSERT *;
SELECT count(*) FROM logs;   -- 2

-- 何度実行しても増えない=冪等
MERGE INTO logs USING newLogs ON logs.uniqueId = newLogs.uniqueId
WHEN NOT MATCHED THEN INSERT *;
SELECT count(*) FROM logs;   -- 2

確認ポイント

  • 3 句の役割: WHEN MATCHED(UPDATE / DELETE)、WHEN NOT MATCHED [BY TARGET]INSERT のみ)、WHEN NOT MATCHED BY SOURCE(UPDATE / DELETE、ターゲット列のみ参照可
  • 同じ種類の句が複数あるとき、最後の句以外は条件が必須
  • 複数一致は失敗する(無条件 DELETE は例外)。CDC では前処理で最新のみ残す
  • insert-only マージは冪等だが、新データセット内部の重複は防げない

課題 8|メダリオン(Bronze→Silver→Gold)を一気通貫で作る

目的: Associate で最も出る「取り込み→洗浄→集計」の型を自分の手で作る。 対応教材: 03-transformation.md 3-6 / 確認問題: Q7

python
from pyspark.sql.functions import col, to_date, sum as _sum, count, avg

# Bronze: 課題 3 で作った bronze_events を使う(生データ+メタデータ列)
# Silver: クレンジング・型キャスト・重複排除・検証
silver = (spark.table(f"{CATALOG}.{SCHEMA}.bronze_events")
    .withColumn("id", col("id").cast("int"))
    .withColumn("event_ts", col("ts").cast("timestamp"))
    .filter(col("id").isNotNull() & col("type").isNotNull())      # 検証
    .dropDuplicates(["id"])                                       # 重複排除
    .select("id", "type", "event_ts", "source_file"))
silver.write.mode("overwrite").saveAsTable(f"{CATALOG}.{SCHEMA}.silver_events")

# Gold: 集計(ディメンションモデル/集計テーブル)
(spark.table(f"{CATALOG}.{SCHEMA}.silver_events")
    .withColumn("event_date", to_date("event_ts"))
    .groupBy("event_date", "type")
    .agg(count("*").alias("events"))
    .write.mode("overwrite").saveAsTable(f"{CATALOG}.{SCHEMA}.gold_events_by_type"))

display(spark.table(f"{CATALOG}.{SCHEMA}.gold_events_by_type"))
sql
-- Gold の集計はマテリアライズドビューでも作れる(SQL ウェアハウス/サーバーレスで)
CREATE OR REPLACE MATERIALIZED VIEW mv_events_by_type AS
SELECT to_date(event_ts) AS event_date, type, count(*) AS events
FROM silver_events GROUP BY 1, 2;

確認ポイント

  • Bronze は生データを元の形式で最小限の検証で保持(多くのフィールドを string / VARIANT / binary で持つのが推奨、メタデータ列を付与)
  • Silver でクレンジング・検証・重複排除・結合を行い、データモデリングを開始(正規化)
  • Gold で集計・ディメンショナルモデリング。Gold は Silver / Bronze より含まれるデータセット数が少ない
  • 取り込みから直接 Silver に書くのは非推奨

課題 9|Lakeflow Jobs で 3 タスクの DAG・cron・通知・Repair run

目的: ジョブ/タスク/トリガーの 3 概念と、失敗タスクだけ再実行できることを体験する。 対応教材: 04-jobs.md 3-1・3-2 / 確認問題: 04-jobs.md Q1〜Q7

手順

  1. ノートブックを 3 つ作る: nb_bronze(課題 3 の取り込み)、nb_silver(課題 8 の Silver)、nb_gold(課題 8 の Gold)
  2. nb_silverわざと失敗させるセルを 1 つ入れる(例: assert False, "intentional failure"
  3. ジョブ & パイプライン → ジョブを作成し、3 タスクを nb_bronze → nb_silver → nb_gold の依存関係でつなぐ
  4. 今すぐ実行 → DAG 上で nb_silver が赤くなり nb_gold がスキップされることを確認
  5. nb_silver の失敗セルを消して [修復して再実行(Repair run)]失敗・スキップしたタスクのみが再実行されることを確認(nb_bronze は再実行されない)
  6. ジョブパラメータ(例 run_date)を追加し、ノートブック側で dbutils.widgets.text("run_date", "") で受け取る
  7. スケジュールとトリガー → トリガーの追加 → スケジュールで cron を設定。[Cron 構文の表示]で Quartz cron を確認し、タイムゾーンに UTC を選ぶ
  8. 通知にメールアドレスを設定し、失敗時に届くことを確認
  9. 確認後、必ずスケジュールを一時停止(PAUSED)する(クォータ節約)

確認ポイント

  • 依存関係は DAG(有向非巡回グラフ) として可視化される
  • Repair run は失敗・スキップしたタスクのサブセットのみを再実行し、成功済みは再実行しない
  • 詳細スケジュールでは Quartz cron 構文が使える。絶対時間で毎時実行したいなら UTC(DST のタイムゾーンだとスキップや 1〜2 時間ずれが起きる)
  • 後続実行の間には最小 10 秒の間隔が強制される
  • 既定の同時実行は 1、超過分はスキップ(キュー待ちではない)
  • Free Edition の制約: 同時実行タスクは最大 5。Slack / Webhook 宛先はワークスペース管理者の通知先設定が必要で、メール通知から試すのが確実

課題 10|宣言型パイプラインと Expectations(warn / drop / fail)

目的: ストリーミングテーブル/マテリアライズドビューと、3 つの違反時アクションの差を[データ品質]タブで見る。 対応教材: 04-jobs.md 3-3・3-4 / 確認問題: Q8・Q9・Q10

新規 → ETL パイプラインでパイプラインを作り、transformations に次の SQL を置きます。

sql
-- Bronze: Auto Loader で増分取り込み(SQL では read_files + STREAM)
CREATE OR REFRESH STREAMING TABLE dlt_bronze
COMMENT 'ボリュームから増分取り込み'
AS SELECT * FROM STREAM read_files(
  '/Volumes/<catalog>/de_sandbox/landing/events',
  format => 'json',
  inferColumnTypes => 'true'
);

-- Silver: 3 つのアクションを並べて違いを見る
CREATE OR REFRESH MATERIALIZED VIEW dlt_silver (
  CONSTRAINT warn_type      EXPECT (type IS NOT NULL),                          -- warn(既定):保持
  CONSTRAINT drop_bad_id    EXPECT (id IS NOT NULL) ON VIOLATION DROP ROW,      -- drop:削除
  CONSTRAINT fail_negative  EXPECT (id >= 0)        ON VIOLATION FAIL UPDATE    -- fail:更新失敗
)
AS SELECT cast(id AS INT) AS id, type, cast(ts AS TIMESTAMP) AS event_ts
FROM STREAM dlt_bronze;

-- Gold
CREATE OR REFRESH MATERIALIZED VIEW dlt_gold AS
SELECT type, count(*) AS n FROM dlt_silver GROUP BY type;

そのうえで:

  1. まず正常データで更新 → [データ品質]タブで warn / drop の合否件数メトリックを確認
  2. id が NULL のレコードを含む JSON を追加 → DROP ROW で削除件数が増えることを確認
  3. id が負のレコードを追加 → FAIL UPDATE で更新が失敗し、メトリックが記録されないことを確認
  4. 隔離(quarantine)パターンも試す(is_quarantined フラグ列+一時テーブル+有効/無効ビュー)

確認ポイント

  • warn=EXPECT(保持・件数記録、既定)/ drop=ON VIOLATION DROP ROW(削除・件数記録)/ fail=ON VIOLATION FAIL UPDATE(更新失敗・アトミックにロールバック・メトリックは記録されない
  • 複数期待値のグループ適用(expect_all 系)は Python 限定。SQL は CONSTRAINT をカンマ区切りで並べるだけ
  • 制約句にカスタム Python 関数・外部呼び出し・他テーブル参照サブクエリは使えない
  • ストリーミングテーブル=追加専用ソースを増分・冪等に取り込む主要型 / マテリアライズドビュー=クエリ結果を実体化
  • Free Edition の制約: パイプライン種別ごとにアクティブ 1 本。課題が終わったらパイプラインを停止/削除する

課題 11|Unity Catalog の権限と USE CATALOG の壁・ボリューム

目的: 「SELECT だけでは読めない」を実感し、3 階層とボリュームの役割を固める。 対応教材: 05-governance.md確認問題: 05-governance.md

sql
-- 3 階層ネームスペースと現在位置
SELECT current_catalog(), current_schema(), current_user();
SHOW SCHEMAS IN <catalog>;
SHOW TABLES IN <catalog>.de_sandbox;

-- 権限の確認(誰に何が付いているか)
SHOW GRANTS ON SCHEMA <catalog>.de_sandbox;
SHOW GRANTS ON TABLE <catalog>.de_sandbox.silver_events;

-- 所有者・詳細メタデータ
DESCRIBE SCHEMA EXTENDED <catalog>.de_sandbox;
DESCRIBE TABLE EXTENDED <catalog>.de_sandbox.silver_events;

-- GRANT の 3 点セット(Free Edition は 1 ユーザーなので、構文と SHOW GRANTS への反映を確認する)
GRANT USE CATALOG ON CATALOG <catalog> TO `account users`;
GRANT USE SCHEMA  ON SCHEMA  <catalog>.de_sandbox TO `account users`;
GRANT SELECT      ON TABLE   <catalog>.de_sandbox.silver_events TO `account users`;
SHOW GRANTS ON TABLE <catalog>.de_sandbox.silver_events;

-- データ検出のための BROWSE
GRANT BROWSE ON CATALOG <catalog> TO `account users`;

-- 取り消し
REVOKE SELECT ON TABLE <catalog>.de_sandbox.silver_events FROM `account users`;
python
# ボリューム = 表形式でないデータのパスベースアクセス(SELECT はできない)
dbutils.fs.ls(f"/Volumes/{CATALOG}/{SCHEMA}/landing")
dbutils.fs.put(f"/Volumes/{CATALOG}/{SCHEMA}/landing/note.txt", "volumes are for non-tabular data", overwrite=True)
print(dbutils.fs.head(f"/Volumes/{CATALOG}/{SCHEMA}/landing/note.txt"))
sql
-- リネージを見る:silver_events → gold_events_by_type の依存が Catalog Explorer の[系列]タブに出る
SELECT type, count(*) FROM <catalog>.de_sandbox.silver_events GROUP BY 1;
-- Catalog Explorer でテーブルを開き[系列]→[系列グラフの表示]、列をクリックして列レベル系列を確認

確認ポイント

  • テーブル読み取りには「テーブルへの SELECT」+「親スキーマへの USE SCHEMA」+「親カタログへの USE CATALOG」の3 点セットが必要。使用特権はそれ自体ではデータにアクセスできない前提条件
  • 権限はカタログ/スキーマから子へ継承される。ただしメタストアレベルの特権は継承されない
  • ボリュームは表形式でないデータ用でパスベースアクセス専用。ボリューム内ファイルをテーブルとして登録することはできない
  • リネージは列レベルまで自動キャプチャされ、表示には最低 BROWSE が必要
  • Free Edition の制約: 外部ロケーション/ストレージ資格情報は作れない(自前クラウドストレージが必要)。マネージド vs 外部テーブルの「DROP 時の挙動」は座学で補う

課題 12|OPTIMIZE / 液体クラスタリング / VACUUM とクエリプロファイル

目的: small files 問題を自分で作って解消し、DESCRIBE HISTORY / DESCRIBE DETAIL で前後を比較する。 対応教材: 06-optimization.md確認問題: 06-optimization.md

python
# 1) わざと small files を作る(20 回に分けて追記)
spark.sql(f"DROP TABLE IF EXISTS {CATALOG}.{SCHEMA}.small_files")
src = spark.table("samples.nyctaxi.trips")
for i in range(20):
    (src.limit(500).write.mode("append" if i else "overwrite")
        .saveAsTable(f"{CATALOG}.{SCHEMA}.small_files"))
sql
-- 2) ファイル数を確認(numFiles が多いことを見る)
DESCRIBE DETAIL small_files;

-- 3) OPTIMIZE(ビンパッキング)→ ファイル数が減る
OPTIMIZE small_files;
DESCRIBE DETAIL small_files;
DESCRIBE HISTORY small_files;   -- operation = OPTIMIZE、operationMetrics を確認

-- 4) もう一度 OPTIMIZE → べき等なので効果なし(numFiles が変わらない)
OPTIMIZE small_files;
DESCRIBE DETAIL small_files;

-- 5) 液体クラスタリング(新規テーブルの推奨手法)
CREATE OR REPLACE TABLE clustered_trips
CLUSTER BY (tpep_pickup_datetime)
AS SELECT * FROM samples.nyctaxi.trips;

DESCRIBE TABLE clustered_trips;          -- clusteringColumns を確認
SHOW TBLPROPERTIES clustered_trips;      -- clusterByAuto など

-- キーを変更したら OPTIMIZE FULL で既存データも再クラスタリング(DBR 16.0 以降)
ALTER TABLE clustered_trips CLUSTER BY (trip_distance);
OPTIMIZE clustered_trips FULL;

-- 自動液体クラスタリング(UC マネージドテーブル)
ALTER TABLE clustered_trips CLUSTER BY AUTO;

-- 6) VACUUM:まず DRY RUN で削除対象を見る(既定保持 7 日なので通常は何も消えない)
VACUUM small_files DRY RUN;
DESCRIBE HISTORY small_files;

クエリプロファイルで性能を見る(Spark UI の代替)

sql
-- SQL ウェアハウスで実行し、クエリ履歴 → 該当クエリ → [クエリプロファイル] を開く
SELECT tpep_pickup_datetime, sum(fare_amount)
FROM samples.nyctaxi.trips
WHERE trip_distance > 2
GROUP BY 1 ORDER BY 2 DESC LIMIT 20;

確認ポイント

  • OPTIMIZE は小さいファイルを結合するビンパッキングで、べき等(2 回目は効果なし)。データの中身は変えないので読み取り結果は前後で同じ
  • ZORDER BY液体クラスタリングのない Delta テーブルの手法。新規テーブルは液体クラスタリングが推奨で、パーティション/ZORDER とは併用不可
  • 液体クラスタリングは既存データを書き換えずにキーを再定義できる。ただし既定では過去のデータに適用されないので、初回有効化/キー変更時は OPTIMIZE FULL
  • クラスタリングキーは最大 4、統計収集列(既定で先頭 32 列)から選ぶ
  • VACUUM の既定保持は 7 日で、実行後は保持期間より古いバージョンへタイムトラベルできなくなる
  • Free Edition の制約: Spark UI は使えないので、スキュー/スピルの読み方(Max > 75%tile の 50% 増、Spill (Memory)/(Disk))は座学で暗記し、実機ではクエリプロファイルでステージ・行数・時間を見る

追加課題|Git フォルダーと Databricks CLI・バンドル

目的: CI/CD の実体(YAML でジョブを宣言し、CLI でデプロイ・実行)を通しで体験する。 対応教材: 07-cicd.md確認問題: 07-cicd.md

A. Git フォルダー

  1. GitHub で空リポジトリを作り、ユーザー設定 → リンクされたアカウント → Git 資格情報の追加で PAT を登録
  2. ワークスペースで 作成 → Git フォルダー からクローン
  3. ノートブックを作って コミット & プッシュ
  4. ブランチを作成して切り替え → 新ブランチに存在しない資産が消えることを確認(元に戻すと新しい ID / URL で再作成される)
  5. GitHub 側で変更してから プルノートブックの状態がリセットされることを確認

B. Databricks CLI とバンドル(ローカル PC で実行)

bash
# 1) CLI を入れる(バンドルには v0.218.0 以降が必要)
winget install Databricks.DatabricksCLI     # Windows
databricks --version

# 2) OAuth U2M でログイン(構成プロファイルが .databrickscfg に保存される)
databricks auth login --host https://<your-workspace-host>

# 3) バンドルを作る
databricks bundle init          # 既定テンプレートを選ぶ

# 4) 検証 → デプロイ → 実行 → 破棄
databricks bundle validate
databricks bundle deploy -t dev
databricks bundle run -t dev <resource_key>
databricks bundle destroy

databricks.yml の最小形を自分で書いてみる:

yaml
bundle:
  name: de_sandbox_bundle

resources:
  jobs:
    hello_job:
      name: hello-job
      tasks:
        - task_key: hello-task
          notebook_task:
            notebook_path: ./src/hello.py

targets:
  dev:
    mode: development
    default: true       # default: true にできるターゲットは 1 つだけ

確認ポイント

  • バンドルのライフサイクルは 作成 → 開発 → 検証(validate)→ デプロイ(deploy)→ 実行(run)→ 破棄(destroy) の 6 段階
  • databricks.ymlルートに 1 つだけの必須メイン構成ファイル。追加は include で参照。default: true は 1 ターゲットのみ
  • mode: development はリソース名に [dev ...] プレフィックス、スケジュール一時停止、同時実行有効、ロック無効。mode: production は権限・ロックが厳格
  • Git フォルダーは対話的開発用、CI/CD と本番デプロイはバンドルが公式推奨
  • 認証設定の評価順序は バンドル設定 → 環境変数 → .databrickscfg プロファイル
  • Free Edition の制約: 1 ワークスペースのみなので dev / prod を別ワークスペースに分ける演習はできない。ターゲットを 2 つ書いて mode の違いだけ確認する

仕上げチェック

全課題を終えたら、05_practice-questions/associate/ の 7 ファイル(各 10 問)を解いて 9 割を目指してください。間違えた問題の対応課題に戻るのが最短ルートです。

最後に片付け(クォータ節約)

python
for q in spark.streams.active:
    q.stop()
sql
DROP SCHEMA IF EXISTS <catalog>.de_sandbox CASCADE;