テーマ切替
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"課題一覧
| # | 課題 | 対応セクション | 所要 |
|---|---|---|---|
| 1 | Delta の ACID・タイムトラベル・トランザクションログを覗く | ① | 20分 |
| 2 | マネージド/外部テーブルとスキーマ適用・進化 | ① ③ | 20分 |
| 3 | Auto Loader で増分取り込み(スキーマ推論・進化・_rescued_data) | ② | 40分 |
| 4 | COPY INTO と冪等性・VALIDATE | ② | 20分 |
| 5 | PySpark DataFrame API と変換/アクション・遅延評価 | ③ | 30分 |
| 6 | UDF 3 種の性能を実測する(組み込み / pandas / 行ごと Python) | ③ | 30分 |
| 7 | MERGE INTO でアップサート・重複排除・SCD Type 2 | ③ | 40分 |
| 8 | メダリオン(Bronze→Silver→Gold)を一気通貫で作る | ③ | 40分 |
| 9 | Lakeflow Jobs で 3 タスクの DAG・cron・通知・Repair run | ④ | 40分 |
| 10 | 宣言型パイプラインと Expectations(warn / drop / fail) | ④ | 40分 |
| 11 | Unity Catalog の権限と USE CATALOG の壁・ボリューム | ⑤ | 30分 |
| 12 | OPTIMIZE / 液体クラスタリング / 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 HISTORYのoperation列に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/selectは narrow、groupBy/join/orderBy/distinctは wide- 運用パイプラインには書き込みアクションのみを置くべき理由(
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
手順
- ノートブックを 3 つ作る:
nb_bronze(課題 3 の取り込み)、nb_silver(課題 8 の Silver)、nb_gold(課題 8 の Gold) nb_silverをわざと失敗させるセルを 1 つ入れる(例:assert False, "intentional failure")- ジョブ & パイプライン → ジョブを作成し、3 タスクを
nb_bronze → nb_silver → nb_goldの依存関係でつなぐ - 今すぐ実行 → DAG 上で
nb_silverが赤くなりnb_goldがスキップされることを確認 nb_silverの失敗セルを消して [修復して再実行(Repair run)] → 失敗・スキップしたタスクのみが再実行されることを確認(nb_bronzeは再実行されない)- ジョブパラメータ(例
run_date)を追加し、ノートブック側でdbutils.widgets.text("run_date", "")で受け取る - スケジュールとトリガー → トリガーの追加 → スケジュールで cron を設定。[Cron 構文の表示]で Quartz cron を確認し、タイムゾーンに UTC を選ぶ
- 通知にメールアドレスを設定し、失敗時に届くことを確認
- 確認後、必ずスケジュールを一時停止(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;そのうえで:
- まず正常データで更新 → [データ品質]タブで warn / drop の合否件数メトリックを確認
idが NULL のレコードを含む JSON を追加 →DROP ROWで削除件数が増えることを確認idが負のレコードを追加 →FAIL UPDATEで更新が失敗し、メトリックが記録されないことを確認- 隔離(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 フォルダー
- GitHub で空リポジトリを作り、ユーザー設定 → リンクされたアカウント → Git 資格情報の追加で PAT を登録
- ワークスペースで 作成 → Git フォルダー からクローン
- ノートブックを作って コミット & プッシュ
- ブランチを作成して切り替え → 新ブランチに存在しない資産が消えることを確認(元に戻すと新しい ID / URL で再作成される)
- 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 destroydatabricks.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;