Skip to content

Free Edition ハンズオン課題|Professional 対策(全15課題)

Databricks 認定データエンジニア Professional の 10 ドメインに対応した、Free Edition で実際に動かせる課題集です。 先に 00-free-edition.md のセットアップと、handson-associate.md の課題(特に 3・7・8)を済ませておくと理解が早いです。 Professional は「シナリオで最適解を選ぶ」試験なので、各課題にわざと壊す/トレードオフを比較する手順を入れています。

前提: <catalog> は自分のワークスペースカタログ名に置き換えます。

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

課題一覧

#課題ドメイン所要
1ウィンドウ関数を総当たりで検算する(RANK / ROWS vs RANGE / 既定フレーム)① 22%40分
2ランキング関数+フレームと RANGE のエラーを踏む15分
3複雑型・explode 系・高階関数を比較する30分
4UDF 5 種を作り分ける(scalar / pandas 4 型 / UDTF / UDAF)+ NULL 落とし穴40分
5液体クラスタリングのデータスキップを数字で確認する② 13%40分
6VACUUM とタイムトラベルのトレードオフを体験する30分
7Expectations の隔離・検証テーブルと「ゲートされない」問題③ 10%40分
8MERGE の句ルールとエラーを全部踏む30分
9df.observe + StreamingQueryListener でストリームを監視する④ 10%40分
10SQL アラートとジョブ通知を設定する30分
11行フィルター・列マスク・動的ビューを 3 通りで実装する⑤ 10%50分
12foreachBatch の at-least-once を壊して冪等化する⑥ ⑦50分
13チェックポイントを壊す/変えると何が起きるか⑦ 7%40分
14リネージ・タグ・BROWSE とガバナンス運用⑧ 7%30分
15SCD Type 1/2・AUTO CDC・CDF を作り比べる⑨ 6%60分
⑩ データ共有・フェデレーションは Free Edition では実質不可 → 座学⑩ 5%

課題 1|ウィンドウ関数を総当たりで検算する

目的: RANK / DENSE_RANK / ROW_NUMBER のタイ挙動、ROWS と RANGE の差ORDER BY 付き集計の既定フレームを実データで確認する。教材の公式例をそのまま再現する。 対応教材: 01-code-development.md 3-3 / 確認問題: 01-code-development.md Q3・Q4・Q6

sql
CREATE OR REPLACE TEMP VIEW employees AS SELECT * FROM VALUES
  ('Fred',  'Engineering', 21000, 23),
  ('Chloe', 'Engineering', 23000, 25),
  ('Tom',   'Engineering', 23000, 28),
  ('Paul',  'Engineering', 29000, 33),
  ('Lisa',  'Sales',       10000, 30),
  ('Alex',  'Sales',       30000, 41),
  ('Evan',  'Sales',       32000, 45)
  AS t(name, dept, salary, age);

-- 1) ランキング関数のタイ挙動(1,2,3 / 1,2,2,4 / 1,2,2,3)
SELECT name, dept, salary,
       ROW_NUMBER() OVER (PARTITION BY dept ORDER BY salary) AS rn,
       RANK()       OVER (PARTITION BY dept ORDER BY salary) AS rnk,
       DENSE_RANK() OVER (PARTITION BY dept ORDER BY salary) AS drnk
FROM employees ORDER BY dept, salary;

-- 2) ROWS と RANGE の決定的な差(Chloe と Tom が同じ 23000)
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' ORDER BY salary;

-- 3) 既定フレームの落とし穴:ORDER BY を付けると累積和になる
SELECT name, dept, salary,
       SUM(salary) OVER (PARTITION BY dept ORDER BY salary)                                   AS default_frame,
       SUM(salary) OVER (PARTITION BY dept)                                                   AS whole_partition,
       SUM(salary) OVER (PARTITION BY dept ORDER BY salary
                         ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING)            AS explicit_all
FROM employees ORDER BY dept, salary;

-- 4) LAG / LEAD の境界とデフォルト値
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_with_default,
       LEAD(salary)       OVER (PARTITION BY dept ORDER BY salary) AS lead_null
FROM employees WHERE dept = 'Sales' ORDER BY salary;

-- 5) 3 行移動平均・現在行から末尾まで・値ベースの範囲
SELECT name, dept, salary,
       ROUND(AVG(salary) OVER (PARTITION BY dept ORDER BY salary
                               ROWS BETWEEN 1 PRECEDING AND 1 FOLLOWING)) AS moving_avg_3,
       SUM(salary) OVER (PARTITION BY dept ORDER BY salary
                         ROWS BETWEEN CURRENT ROW AND UNBOUNDED FOLLOWING) AS to_end,
       SUM(salary) OVER (PARTITION BY dept ORDER BY salary
                         RANGE BETWEEN 5000 PRECEDING AND 5000 FOLLOWING)  AS pm5000
FROM employees ORDER BY dept, salary;

確認ポイント

  • Engineering の rows_total21000 / 44000 / 67000 / 96000range_total21000 / 67000 / 67000 / 96000タイの先頭行だけ差が出る
  • default_framewhole_partition と一致せず、explicit_all と一致する — ORDER BY を付けたら累積和
  • LEAD はデフォルト未指定なら NULLLEAD(col, 1, 0) なら 0
  • PySpark でも同じことを書けるか(Window.partitionBy().orderBy().rowsBetween(...)

課題 2|ランキング関数+フレームと RANGE のエラーを踏む

目的: エラーコードを実際に見て覚える。試験ではエラー名が選択肢に出る。 対応教材: 同 3-3 / 確認問題: Q5

sql
-- 1) ランキング関数にフレームを付ける → WINDOW_FUNCTION_AND_FRAME_MISMATCH
SELECT name, RANK() OVER (PARTITION BY dept ORDER BY salary
                          ROWS BETWEEN 1 PRECEDING AND CURRENT ROW) AS r
FROM employees;

-- 2) RANGE フレームに ORDER BY がない → DATATYPE_MISMATCH.RANGE_FRAME_WITHOUT_ORDER
SELECT name, SUM(salary) OVER (PARTITION BY dept
                               RANGE BETWEEN 5000 PRECEDING AND CURRENT ROW) AS s
FROM employees;

-- 3) RANGE フレームで ORDER BY が複数式 → DATATYPE_MISMATCH.RANGE_FRAME_MULTI_ORDER
SELECT name, SUM(salary) OVER (PARTITION BY dept ORDER BY salary, age
                               RANGE BETWEEN 5000 PRECEDING AND CURRENT ROW) AS s
FROM employees;

確認ポイント

  • ランキング関数は ORDER BY 必須・フレーム句は禁止
  • RANGE フレームは ORDER BY 必須・式は 1 つのみ
  • 3 つのエラーコードを見分けられるか

課題 3|複雑型・explode 系・高階関数を比較する

目的: explodeexplode_outer の差、そして「高階関数ならシャッフルなしで配列を直接操作できる」を実行計画で確認する。 対応教材: 同 3-4 / 確認問題: Q10

sql
CREATE OR REPLACE TEMP VIEW orders_nested AS SELECT * FROM VALUES
  (1, array(10, 20, 30), named_struct('name','Alice','addr', named_struct('city','Tokyo')), map('k1','v1')),
  (2, array(),           named_struct('name','Bob',  'addr', named_struct('city','Osaka')), map('k2','v2')),
  (3, cast(NULL AS array<int>), named_struct('name','Carol','addr', named_struct('city','Kyoto')), map())
  AS t(id, items, customer, tags);

-- 1) explode は空/NULL の行が消える、explode_outer は残る、posexplode は位置も返す
SELECT id, explode(items)       AS item FROM orders_nested;   -- id=2,3 が消える
SELECT id, explode_outer(items) AS item FROM orders_nested;   -- id=2,3 が NULL 行で残る
SELECT id, posexplode(items)    AS (pos, item) FROM orders_nested;

-- 2) struct のドットアクセス(ネストも可)
SELECT id, customer.name, customer.addr.city FROM orders_nested;

-- 3) 高階関数:配列を直接操作(explode → 集計 → collect_list の往復が不要)
SELECT id,
       transform(items, x -> x * 2)                                        AS doubled,
       filter(items, x -> x > 15)                                          AS gt15,
       exists(items, x -> x > 25)                                          AS has_big,
       forall(items, x -> x > 0)                                           AS all_positive,
       aggregate(items, 0, (acc, x) -> acc + x)                            AS total,
       aggregate(filter(transform(items, x -> x * x), y -> y > 200), 0, (acc, v) -> acc + v) AS nested
FROM orders_nested;
python
# 4) 実行計画を比較する:高階関数はシャッフルなし、explode + groupBy はシャッフルあり
hof = spark.sql("SELECT id, aggregate(items, 0, (acc, x) -> acc + x) AS total FROM orders_nested")
exp = spark.sql("SELECT id, sum(item) AS total FROM (SELECT id, explode(items) AS item FROM orders_nested) GROUP BY id")

hof.explain("formatted")   # Exchange が出ない
exp.explain("formatted")   # Exchange(シャッフル)が出る

確認ポイント

  • explode空/NULL の要素で行が消えるexplode_outer は NULL 行を保持、posexplode は位置も返す、inline は構造体配列を複数列に展開
  • 高階関数は explode → 集計 → collect_list のシャッフルを伴う往復を避けられる(実行計画の Exchange の有無で確認)
  • ラムダ構文 x -> expr と 2 引数 (acc, x) -> expr を読める/書けるか

課題 4|UDF 5 種を作り分ける + NULL 落とし穴

目的: 入出力から UDF 種別を判別できるようにする。UC 管理とセッションスコープの差も見る。 対応教材: 同 3-5 / 確認問題: Q7・Q8・Q9

python
import pandas as pd
from typing import Iterator, Tuple
from pyspark.sql import Window
from pyspark.sql.functions import udf, pandas_udf, udtf, lit, col

df = spark.createDataFrame([(1, 70.0, 1.75), (2, 80.0, 1.80), (1, 60.0, 1.65)], "id INT, w DOUBLE, h DOUBLE")

# ① スカラー UDF:1 行 → 1 値
@udf("double")
def bmi_scalar(w, h):
    return w / (h ** 2)

# ② pandas UDF: Series → Series(select / withColumn 用)
@pandas_udf("double")
def bmi_s2s(w: pd.Series, h: pd.Series) -> pd.Series:
    return w / (h ** 2)

# ③ pandas UDF: Iterator[複数 Series] → Iterator[Series](状態初期化ができる)
@pandas_udf("double")
def bmi_iter(it: Iterator[Tuple[pd.Series, pd.Series]]) -> Iterator[pd.Series]:
    # ここで重い初期化を 1 回だけ行える(ML モデルのロードなど)
    for w, h in it:
        yield w / (h ** 2)

# ④ pandas UDF: Series → Scalar(=UDAF。集計・Window で使う。セッションスコープ限定)
@pandas_udf("double")
def mean_udaf(v: pd.Series) -> float:
    return v.mean()

# ⑤ UDTF:1 行 → 複数行(複数列)
@udtf(returnType="sum: int, diff: int")
class GetSumDiff:
    def eval(self, x: int, y: int):
        yield x + y, x - y

df.select("*", bmi_scalar("w", "h").alias("bmi1"),
               bmi_s2s("w", "h").alias("bmi2"),
               bmi_iter("w", "h").alias("bmi3")).show()

df.groupBy("id").agg(mean_udaf("w").alias("avg_w")).show()

w = Window.partitionBy("id").rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing)
df.withColumn("mean_w", mean_udaf("w").over(w)).show()

GetSumDiff(lit(10), lit(3)).show()
python
# NULL 評価順序の落とし穴を実際に踏む
spark.sql("CREATE OR REPLACE TEMP VIEW test1 AS SELECT * FROM VALUES ('abc'), (NULL) AS t(s)")

spark.udf.register("strlen", lambda s: len(s), "int")            # NULL 非対応
try:
    spark.sql("SELECT s FROM test1 WHERE s IS NOT NULL AND strlen(s) > 1").show()
except Exception as e:
    print("短絡が保証されないため失敗:", str(e)[:160])

# 対策 1: UDF 自体を NULL 対応にする
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").show()

# 対策 2: IF / CASE WHEN で条件分岐内で呼ぶ
spark.sql("SELECT s FROM test1 WHERE if(s IS NOT NULL, strlen(s), NULL) > 1").show()
sql
-- UC 管理 UDF(永続化・共有可)とセッションスコープの違い
CREATE OR REPLACE FUNCTION <catalog>.de_sandbox.get_name_length(name STRING)
RETURNS INT RETURN LENGTH(name);

SELECT <catalog>.de_sandbox.get_name_length('databricks');
SHOW FUNCTIONS IN <catalog>.de_sandbox;      -- カタログエクスプローラーからも検出できる
GRANT EXECUTE ON FUNCTION <catalog>.de_sandbox.get_name_length TO `account users`;

確認ポイント

  • 入出力から判別: 1 行→1 値=スカラー / 複数行→1 値=UDAF / 1 行→複数行=UDTF / Series→Series・Scalar=pandas UDF
  • pandas UDF の 4 型を型ヒントで区別できるか
  • UDAF はセッションスコープ限定、UC 管理 UDF は SQL / Python / Scala / Java 対応で永続化・共有・ガバナンス可
  • UDF の NULL 落とし穴(短絡保証がない)と 2 つの対策
  • Free Edition の制約: UDF は外部ネットワーク不可・メモリ 1 GB 上限(教材の記述をそのまま体験できる)。Scala UDF は非サポート

課題 5|液体クラスタリングのデータスキップを数字で確認する

目的: 「クラスタリングキーでフィルターすると読むファイルが減る」をクエリプロファイルの数字で見る。 対応教材: 02-cost-performance.md 3-2 / 確認問題: 02-cost-performance.md Q1〜Q3

sql
-- 比較用に 2 つ作る:クラスタリングなし / あり
CREATE OR REPLACE TABLE trips_plain AS
SELECT * FROM samples.nyctaxi.trips;

CREATE OR REPLACE TABLE trips_clustered
CLUSTER BY (tpep_pickup_datetime, trip_distance) AS
SELECT * FROM samples.nyctaxi.trips;

-- 初回有効化なので既存データも再クラスタリングする(DBR 16.4 LTS 以降)
OPTIMIZE trips_clustered FULL;

DESCRIBE DETAIL trips_plain;
DESCRIBE DETAIL trips_clustered;
DESCRIBE TABLE trips_clustered;          -- clusteringColumns を確認
SHOW TBLPROPERTIES trips_clustered;
sql
-- SQL ウェアハウスで両方を実行し、クエリプロファイルの「読み取りファイル数 / 読み取りバイト数」を比較
SELECT count(*) FROM trips_plain
WHERE tpep_pickup_datetime BETWEEN '2016-02-01' AND '2016-02-03' AND trip_distance > 5;

SELECT count(*) FROM trips_clustered
WHERE tpep_pickup_datetime BETWEEN '2016-02-01' AND '2016-02-03' AND trip_distance > 5;
sql
-- キー変更 → 増分 OPTIMIZE では既存データが書き換わらないことを確認
ALTER TABLE trips_clustered CLUSTER BY (trip_distance);
OPTIMIZE trips_clustered;         -- 増分:既存の一致しないファイルは書き換えない
DESCRIBE HISTORY trips_clustered;

OPTIMIZE trips_clustered FULL;    -- 強制再クラスタリング
DESCRIBE HISTORY trips_clustered; -- operationMetrics の numFilesAdded/Removed を比較

-- 自動液体クラスタリング(UC マネージドテーブル+予測的最適化が前提)
ALTER TABLE trips_clustered CLUSTER BY AUTO;
SHOW TBLPROPERTIES trips_clustered;   -- clusterByAuto = true

-- パーティション/ZORDER とは併用できないことを確認(エラーになる)
CREATE OR REPLACE TABLE bad_combo (a INT, b STRING)
PARTITIONED BY (b) CLUSTER BY (a);

確認ポイント

  • キーは最大 4、**統計収集列(既定で先頭 32 列)**から選ぶ。10 TB 未満ではキーが多いと単一列フィルターが遅くなりうる
  • OPTIMIZE は増分。初回有効化/キー変更時は OPTIMIZE FULL
  • パーティション/ZORDER と併用不可
  • 書き込み時クラスタリングのしきい値(UC マネージド: 1 キー 64MB / 2 キー 256MB / 3 キー 512MB / 4 キー 1GB。その他 Delta はこの 4 倍)は座学で暗記
  • Free Edition の制約: サーバーレスなのでノード種別(コンピュート最適化など)は選べない。「OPTIMIZE は CPU 集約なのでコンピュート最適化インスタンス推奨」は座学

課題 6|VACUUM とタイムトラベルのトレードオフを体験する

目的: 「VACUUM するとタイムトラベルできなくなる」を安全な範囲で確認する。 対応教材: 同 3-2 / 確認問題: Q5

sql
CREATE OR REPLACE TABLE t_vacuum AS SELECT * FROM samples.nyctaxi.trips LIMIT 1000;
UPDATE t_vacuum SET fare_amount = fare_amount * 1.1;   -- v1:古いファイルが未参照になる
DELETE FROM t_vacuum WHERE fare_amount > 100;          -- v2

DESCRIBE HISTORY t_vacuum;
SELECT count(*) FROM t_vacuum VERSION AS OF 0;   -- まだ読める

-- 1) DRY RUN で削除候補を見る(既定保持 7 日なので候補は出ないはず)
VACUUM t_vacuum DRY RUN;

-- 2) 保持期間を 0 にして実行 → 安全チェックに弾かれることを確認(重要)
VACUUM t_vacuum RETAIN 0 HOURS;
sql
-- 3) 安全チェックを外して実行(学習用テーブルでのみ!本番では絶対にやらない)
SET spark.databricks.delta.retentionDurationCheck.enabled = false;
VACUUM t_vacuum RETAIN 0 HOURS DRY RUN;   -- 削除候補が出る
VACUUM t_vacuum RETAIN 0 HOURS;

-- 4) 古いバージョンが読めなくなったことを確認
SELECT count(*) FROM t_vacuum VERSION AS OF 0;   -- ファイルが無くエラーになる

SET spark.databricks.delta.retentionDurationCheck.enabled = true;   -- 必ず戻す

確認ポイント

  • 既定保持は 7 日で、Databricks は 7 日以上を強く推奨。理由は数日実行されるジョブの未コミットファイルが完了前に削除される恐れがあるから
  • 安全チェックは spark.databricks.delta.retentionDurationCheck.enabled(Iceberg は spark.databricks.iceberg.retentionDurationCheck.enabled
  • VACUUM 実行後は保持期間より古いバージョンにタイムトラベルできない
  • ログファイルは既定保持 30 日で**VACUUM の管理外**
  • 長期タイムトラベルが必要なら、予測的最適化を有効化する前に delta.deletedFileRetentionDuration を設定する
  • LITE / FULL モードは DBR 16.4 LTS 以降のパブリックプレビュー。LITE はログ保持しきい値(既定 30 日)内に成功した VACUUM が 1 回以上必要

課題 7|Expectations の隔離・検証テーブルと「ゲートされない」問題

目的: 「検証テーブルはオーケストレーションではない」を実機で確認する。Professional 特有の落とし穴。 対応教材: 03-transformation-quality.md 3-2・4-3・4-4 / 確認問題: 03-transformation-quality.md Q4・Q5

Lakeflow パイプラインを作り、transformations に Python ファイルとして置きます。

python
from pyspark import pipelines as dp
from pyspark.sql.functions import expr, col

# --- 隔離(quarantine)パターン:データを捨てずに有効/無効を分離 ---
rules = {
    "valid_id":   "(id IS NOT NULL)",
    "valid_type": "(type IS NOT NULL)",
}
quarantine_rules = "NOT({0})".format(" AND ".join(rules.values()))

@dp.view
def raw_events():
    return spark.readStream.table(f"<catalog>.de_sandbox.bronze_events")

@dp.table(temporary=True, partition_cols=["is_quarantined"])
@dp.expect_all(rules)                       # warn なのでメトリックは取れる
def events_quarantine():
    return (spark.readStream.table("raw_events")
            .withColumn("is_quarantined", expr(quarantine_rules)))

@dp.view
def valid_events():
    return spark.read.table("events_quarantine").filter("is_quarantined = false")

@dp.view
def invalid_events():
    return spark.read.table("events_quarantine").filter("is_quarantined = true")

# --- 検証テーブル:行数一致を expect_or_fail でゲートしたい(が、下流は待たない)---
@dp.materialized_view(name="count_verification")
@dp.expect_or_fail("no_rows_dropped", "a_count == b_count")
def validate_row_counts():
    return spark.sql("""
      SELECT * FROM
        (SELECT COUNT(*) AS a_count FROM LIVE.valid_events),
        (SELECT COUNT(*) AS b_count FROM LIVE.events_quarantine)
    """)

# --- 主キー一意性の検証 ---
@dp.materialized_view(name="pk_verification")
@dp.expect_or_fail("unique_pk", "num_entries = 1")
def validate_pk():
    return (spark.read.table("events_quarantine").groupBy("id").count()
            .withColumnRenamed("count", "num_entries"))

# --- これが「ゲートされない」ダウンストリーム ---
@dp.materialized_view
def gold_by_type():
    return spark.sql("SELECT type, count(*) AS n FROM LIVE.valid_events GROUP BY type")

やること

  1. 正常データで更新 → [データ品質]タブで warn のメトリックを確認
  2. id が NULL のレコードを追加して更新 → count_verificationfail するのを確認
  3. それでも gold_by_type が更新されていることを確認 → これが「Expectations はオーケストレーションではない」
  4. 正しい設計に直す: 検証を別パイプラインに分け、ジョブで 検証パイプライン → 本体パイプライン のタスク依存にする

確認ポイント

  • 検証テーブルを他データセットから読んでも、そのデータセットは検証結果を待たない。止めたいなら別パイプライン+ジョブのタスク依存
  • expect_or_failメトリックが記録されない(更新が失敗するため)
  • fail 時の挙動はモード依存: トリガー型は当該フローのみロールバック、連続型はフロー+全依存フローが停止
  • expect_all 系は Python 限定
  • Free Edition の制約: パイプラインはアクティブ 1 本なので、「別パイプラインに分割」は 2 本目を作って片方を停止しながら確認する

課題 8|MERGE の句ルールとエラーを全部踏む

目的: エラーメッセージを実際に見る。Professional の MERGE は句ルールが問われる。 対応教材: 同 3-4 / 確認問題: Q6・Q7・Q8

sql
CREATE OR REPLACE TABLE tgt (key INT, value STRING, deleted_count INT, marked_for_deletion BOOLEAN, updated_at TIMESTAMP);
INSERT INTO tgt VALUES (1,'a',0,false,'2026-07-01'), (2,'b',0,true,'2026-07-01'), (3,'c',0,false,'2026-07-01');

CREATE OR REPLACE TEMP VIEW src AS SELECT * FROM VALUES
  (1,'A',TIMESTAMP'2026-07-02'), (4,'D',TIMESTAMP'2026-07-02') AS t(key, value, updated_at);

-- 1) 同じ種類の句が複数あるとき、最後以外に条件がないとエラー
--    → NON_LAST_MATCHED_CLAUSE_OMIT_CONDITION
MERGE INTO tgt USING src ON tgt.key = src.key
WHEN MATCHED THEN UPDATE SET tgt.value = src.value
WHEN MATCHED AND tgt.marked_for_deletion THEN DELETE;

-- 2) 正しい順序(条件付きを先に、無条件を最後に)
MERGE INTO tgt USING src ON tgt.key = src.key
WHEN MATCHED AND tgt.marked_for_deletion THEN DELETE
WHEN MATCHED THEN UPDATE SET tgt.value = src.value, tgt.updated_at = src.updated_at
WHEN NOT MATCHED THEN INSERT (key, value, deleted_count, marked_for_deletion, updated_at)
                      VALUES (src.key, src.value, 0, false, src.updated_at);
SELECT * FROM tgt ORDER BY key;

-- 3) NOT MATCHED BY SOURCE はターゲット列のみ参照可(ソース列を書くとエラー)
MERGE INTO tgt USING src ON tgt.key = src.key
WHEN NOT MATCHED BY SOURCE AND src.value IS NULL THEN DELETE;   -- ← エラー

-- 4) 正しい書き方(ターゲット列 / リテラル / ターゲット列演算)
MERGE INTO tgt USING src ON tgt.key = src.key
WHEN NOT MATCHED BY SOURCE AND tgt.updated_at < TIMESTAMP'2026-07-02'
  THEN UPDATE SET tgt.deleted_count = tgt.deleted_count + 1;
SELECT * FROM tgt ORDER BY key;

-- 5) UPDATE SET * は列名一致が前提(列名が違うと分析エラー)
CREATE OR REPLACE TEMP VIEW src_diff AS SELECT * FROM VALUES (1,'X') AS t(key, val);   -- value ではなく val
MERGE INTO tgt USING src_diff ON tgt.key = src_diff.key WHEN MATCHED THEN UPDATE SET *;

確認ポイント

  • 各種類の句は複数指定でき指定順に評価最後の句以外は条件が必須
  • WHEN NOT MATCHED BY SOURCE はターゲット列のみ参照可、条件なしだとテーブル全書き換えになる
  • UPDATE SET * / INSERT *列名一致が前提
  • MERGE WITH SCHEMA EVOLUTION(DBR 15.2 以上)でターゲットをソースに合わせて自動進化できる

課題 9|df.observe + StreamingQueryListener でストリームを監視する

目的: 入力レート/処理レート/バックログ/カスタムメトリックをイベントとして取り出す対応教材: 04-monitoring-alerting.md 3-3 / 確認問題: 04-monitoring-alerting.md Q7・Q8

python
from pyspark.sql.streaming import StreamingQueryListener
from pyspark.sql.functions import count, lit, col, when

class MyListener(StreamingQueryListener):
    def onQueryStarted(self, event):
        print(f"[started] name={event.name} id={event.id}")

    def onQueryProgress(self, event):
        p = event.progress
        print(f"[progress] batchId={p.batchId} numInputRows={p.numInputRows} "
              f"in/s={p.inputRowsPerSecond} proc/s={p.processedRowsPerSecond}")
        # durationMs でボトルネック段階を見る
        print("           durationMs:", dict(p.durationMs))
        # バックログ(Delta / Auto Loader)
        for s in p.sources:
            m = dict(s.metrics or {})
            for k in ("numFilesOutstanding", "numBytesOutstanding"):
                if k in m:
                    print(f"           backlog {k}={m[k]}")
        # observe で定義したカスタムメトリック
        om = p.observedMetrics.get("dq")
        if om is not None and om.cnt:
            ratio = om.malformed / om.cnt
            print(f"           malformed ratio={ratio:.2%}")
            if ratio > 0.3:
                print("           *** ALERT: too many malformed records ***")

    def onQueryIdle(self, event):
        pass

    def onQueryTerminated(self, event):
        print(f"[terminated] id={event.id} exception={event.exception}")

listener = MyListener()
spark.streams.addListener(listener)
python
# observe でデータ品質メトリックを定義してストリームを流す
src = (spark.readStream.format("cloudFiles")
    .option("cloudFiles.format", "json")
    .option("cloudFiles.schemaLocation", f"{VOL}/_ckpt/obs")
    .load(f"{VOL}/events"))

observed = src.observe(
    "dq",
    count(lit(1)).alias("cnt"),
    count(when(col("id").isNull(), lit(1))).alias("malformed"),
)

q = (observed.writeStream
     .queryName("obs_ingest")                       # Spark UI/リスナーで識別しやすくする
     .option("checkpointLocation", f"{VOL}/_ckpt/obs")
     .trigger(availableNow=True)
     .toTable(f"{CATALOG}.{SCHEMA}.obs_events"))
q.awaitTermination()

# 最終の進捗を確認
print(q.lastProgress)

spark.streams.removeListener(listener)

確認ポイント

  • コールバックは onQueryStarted / onQueryProgress / onQueryIdle / onQueryTerminatedonQueryProgressマイクロバッチの最後にのみ配信される
  • inputRowsPerSecond > processedRowsPerSecond が続く=バックプレッシャー
  • バックログは Delta / Auto Loader=numFilesOutstanding / numBytesOutstandingKafka=*OffsetsBehindLatest
  • durationMs の段階(triggerExecution / addBatch / getBatch / latestOffset / queryPlanning / walCommit / commitOffsets)でボトルネックを切り分ける
  • リスナーの処理は軽くし、Kafka のような高速シンクへ書くのが推奨
  • Free Edition の制約: Spark UI の Streaming タブは使えないので、リスナーと lastProgress が実質唯一の観測手段。reservoirId は UC の tableId ではなく DESCRIBE DETAILid に対応することも確認しておく

課題 10|SQL アラートとジョブ通知を設定する

目的: プッシュ型の監視(条件成立 → 通知)を作る。 対応教材: 同 3-4 / 確認問題: Q9・Q10

A. SQL アラート

  1. SQL エディターでしきい値監視クエリを書く
    sql
    SELECT to_date(tpep_pickup_datetime) AS date, SUM(fare_amount) AS amount
    FROM samples.nyctaxi.trips
    GROUP BY ALL ORDER BY 1 DESC;
  2. アラートを作成し、条件を「値=SUM / 列=amount / 演算子=> / しきい値=静的な値」で設定
  3. [テスト条件] でプレビューし、TRIGGERED になるしきい値を探す
  4. スケジュールを設定([Cron 構文]で Quartz cron を確認)
  5. 通知先に自分のメールを設定し、OK に戻るまでの再通知間隔も設定してみる
  6. 詳細設定で「OK 時に通知」「空の結果状態」「通知テンプレート」を確認({{ALERT_STATUS}} {{QUERY_RESULT_TABLE}} などの変数)
  7. 確認後、スケジュールを止める

B. ジョブ通知

  1. 課題 9 のジョブに 通知を追加し、「開始 / 成功 / 失敗 / 期間がしきい値超過」を設定
  2. 「期間がしきい値超過」を試すには予想期間(タイムアウト)を設定する必要があることを確認
  3. わざと失敗するタスクを入れ、タスクを再試行させる設定にして「再試行中はジョブレベル通知が飛ばない」ことを確認 → タスク通知を追加して届くようにする
  4. 「スキップされた実行の通知をミュート」「最後の再試行まで通知をミュート」などのフィルタを試す

確認ポイント

  • SQL アラートの状態は OK / TRIGGERED / ERRORUNKNOWN は廃止)。パラメーター付きクエリ非対応既存の保存済みクエリは再利用不可
  • アラート実行はワークスペースのアクティブ実行クォータ(既定 250)を他ジョブ種と共有し、埋まっているとスキップされる → スケジュールをずらす
  • ジョブ通知: 「失敗を伴う成功」は成功扱い再試行中のタスクはジョブレベル通知が飛ばない(タスク通知を使う)、イベント種別ごと最大 3 システム宛先
  • ストリーミングバックログ通知は 10 分間の平均で判定し、送信後 30 分待つ。ジョブ内で awaitTermination() を使わない
  • Free Edition の制約: Slack / Teams / PagerDuty / Webhook 宛先はワークスペース管理者が通知先を作成する必要がある。メール通知で代替可

課題 11|行フィルター・列マスク・動的ビューを 3 通りで実装する

目的: 3 メカニズムを同じ要件に対して実装し、使い分けを体で覚える。 対応教材: 05-security-compliance.md確認問題: 05-security-compliance.md

sql
CREATE OR REPLACE TABLE customers (
  id INT, name STRING, ssn STRING, region STRING, address STRING, country STRING
);
INSERT INTO customers VALUES
  (1,'Alice','123-45-6789','US','1 Main St','us'),
  (2,'Bob',  '987-65-4321','EU','2 High St','eu'),
  (3,'Carol','555-00-1111','US','3 Oak Ave','us');

A. テーブルレベルの行フィルター

sql
CREATE OR REPLACE FUNCTION us_filter(region STRING)
RETURN IF(is_account_group_member('admins'), true, region = 'US');

ALTER TABLE customers SET ROW FILTER us_filter ON (region);
SELECT * FROM customers;                       -- admins でなければ US のみ

-- 個数制約:1 テーブルに行フィルターは 1 つだけ
DESCRIBE TABLE EXTENDED customers;             -- Row Filter を確認

-- 解除(関数を消す前に必ずこれをやる)
ALTER TABLE customers DROP ROW FILTER;
DROP FUNCTION us_filter;

B. テーブルレベルの列マスク(USING COLUMNS つき)

sql
CREATE OR REPLACE FUNCTION ssn_mask(ssn STRING)
RETURN CASE WHEN is_account_group_member('HumanResourceDept') THEN ssn ELSE '***-**-****' END;

ALTER TABLE customers ALTER COLUMN ssn SET MASK ssn_mask;
SELECT id, name, ssn FROM customers;

-- 他の列を条件に使うマスク
CREATE OR REPLACE FUNCTION mask_address_by_country(
  address STRING, country STRING, group_suffix STRING DEFAULT '_address_viewers')
RETURN IF(is_account_group_member(country || group_suffix), address, 'REDACTED');

ALTER TABLE customers ALTER COLUMN address
  SET MASK mask_address_by_country USING COLUMNS (country, '_address_viewers');
SELECT id, address, country FROM customers;

-- Python UDF は SQL でラップしないと [ROUTINE_NOT_FOUND] になる(実際に踏む)
CREATE OR REPLACE FUNCTION email_mask_python(email STRING) RETURNS STRING
LANGUAGE PYTHON AS $$
import re
return re.sub(r'^[^@]+', lambda m: '*' * len(m.group()), email)
$$;
ALTER TABLE customers ALTER COLUMN name SET MASK email_mask_python;   -- ← エラーになる

CREATE OR REPLACE FUNCTION email_mask_sql(email STRING) RETURN email_mask_python(email);
ALTER TABLE customers ALTER COLUMN name SET MASK email_mask_sql;      -- こちらは通る

ALTER TABLE customers ALTER COLUMN ssn  DROP MASK;
ALTER TABLE customers ALTER COLUMN name DROP MASK;

C. 動的ビュー

sql
CREATE OR REPLACE VIEW customers_redacted AS
SELECT id, name, region,
       CASE WHEN is_account_group_member('auditors') THEN ssn ELSE 'REDACTED' END AS ssn,
       CASE WHEN is_account_group_member('auditors') THEN address
            ELSE regexp_extract(address, '^.*\\s(.*)$', 1) END AS address_part
FROM customers
WHERE CASE WHEN is_account_group_member('managers') THEN TRUE ELSE region = 'US' END;

SELECT * FROM customers_redacted;
SELECT current_user(), session_user();

D. 制限を実際に確認する

sql
-- ビューには行フィルター/列マスクを適用できない
ALTER VIEW customers_redacted SET ROW FILTER us_filter ON (region);   -- ← エラー

-- 行フィルター/列マスク付きテーブルはタイムトラベル・クローン不可
ALTER TABLE customers ALTER COLUMN ssn SET MASK ssn_mask;
SELECT * FROM customers VERSION AS OF 0;                    -- ← エラー
CREATE TABLE customers_clone SHALLOW CLONE customers;       -- ← エラー
ALTER TABLE customers ALTER COLUMN ssn DROP MASK;

-- ANSI モードと型不一致で誤結果になることを確認(重要)
SET spark.sql.ansi.enabled = false;
CREATE OR REPLACE FUNCTION bad_filter(region INT) RETURNS BOOLEAN RETURN region IS NULL;
ALTER TABLE customers SET ROW FILTER bad_filter ON (region);
SELECT * FROM customers;      -- STRING が INT にキャストできず NULL → 全行返ってしまう
ALTER TABLE customers DROP ROW FILTER;
SET spark.sql.ansi.enabled = true;   -- 必ず戻す(Databricks 推奨は有効)

確認ポイント

  • 行フィルターは 1 テーブルに 1 つ、列マスクは 1 列に 1 つ。列マスクの戻り値型は列型と一致/キャスト可能
  • 関数削除は必ず DROP ROW FILTER / DROP MASK の後に DROP FUNCTION
  • ビューには適用できないタイムトラベル・クローン不可Delta Lake API 不可生成列が参照する列にマスク不可
  • ANSI モード無効だとサイレント NULL 化で誤結果(全行返る行フィルター)。Databricks は ANSI 有効を推奨
  • Python UDF は SQL ラッパー経由でのみ列マスクに使える
  • REPLACE TABLE してもフィルター/マスクは保持される(誤消去防止)ことも試してみる
  • Free Edition の制約: グループを作れないため is_account_group_member('...') は基本 false を返す前提で挙動を見る。ABAC の管理タグはアカウントレベル機能なので制限される可能性が高いCREATE POLICY の構文は座学(確認問題 Q7)で固める

課題 12|foreachBatch の at-least-once を壊して冪等化する

目的: 「foreachBatch は at-least-once」を実際に重複させて確認し、txnAppId / txnVersion で冪等化する。Professional の核心。 対応教材: 07-ingestion.md 3-4 / 06-debugging-deploy.md 3-2 / 確認問題: 07-ingestion.md Q9

python
# 準備:ソーステーブル
spark.sql(f"CREATE OR REPLACE TABLE {CATALOG}.{SCHEMA}.raw_orders AS "
          f"SELECT id, id * 100 AS amount FROM range(1000)")
spark.sql(f"DROP TABLE IF EXISTS {CATALOG}.{SCHEMA}.orders_naive")
spark.sql(f"DROP TABLE IF EXISTS {CATALOG}.{SCHEMA}.orders_idem")
python
# --- ① 冪等でない書き込み:チェックポイントを消して再実行すると重複する ---
def naive(batch_df, batch_id):
    batch_df.write.format("delta").mode("append").saveAsTable(f"{CATALOG}.{SCHEMA}.orders_naive")

ck_naive = f"{VOL}/_ckpt/naive"
q = (spark.readStream.table(f"{CATALOG}.{SCHEMA}.raw_orders")
     .writeStream.foreachBatch(naive)
     .option("checkpointLocation", ck_naive)
     .trigger(availableNow=True).start())
q.awaitTermination()
print("1回目:", spark.table(f"{CATALOG}.{SCHEMA}.orders_naive").count())

dbutils.fs.rm(ck_naive, True)          # 障害でチェックポイントを失った状況を再現
q = (spark.readStream.table(f"{CATALOG}.{SCHEMA}.raw_orders")
     .writeStream.foreachBatch(naive)
     .option("checkpointLocation", ck_naive)
     .trigger(availableNow=True).start())
q.awaitTermination()
print("2回目(重複する):", spark.table(f"{CATALOG}.{SCHEMA}.orders_naive").count())   # 2000
python
# --- ② txnAppId + txnVersion で冪等化 ---
ck_idem = f"{VOL}/_ckpt/idem"

def run_idem(app_id):
    """app_id を引数で受け取り、foreachBatch のクロージャに閉じ込める"""
    def idempotent(batch_df, batch_id):
        if batch_df.isEmpty():                       # 空バッチは渡されうる
            return
        (batch_df.write.format("delta").mode("append")
            .option("txnAppId", app_id)              # 一意のアプリ ID
            .option("txnVersion", batch_id)          # batchId にバインド
            .saveAsTable(f"{CATALOG}.{SCHEMA}.orders_idem"))

    q = (spark.readStream.table(f"{CATALOG}.{SCHEMA}.raw_orders")
         .writeStream.foreachBatch(idempotent)
         .option("checkpointLocation", ck_idem)
         .trigger(availableNow=True).start())
    q.awaitTermination()

run_idem("orders-ingest-v1")
print("1回目:", spark.table(f"{CATALOG}.{SCHEMA}.orders_idem").count())      # 1000

dbutils.fs.rm(ck_idem, True)                 # チェックポイント喪失を再現
run_idem("orders-ingest-v1")                 # 同じ txnAppId → 既知バッチはスキップされる
print("2回目(重複しない):", spark.table(f"{CATALOG}.{SCHEMA}.orders_idem").count())   # 1000

# --- ③ 注意点:新チェックポイントでは別の txnAppId が必要 ---
# 新チェックポイントは batchId 0 から始まるため、同じ txnAppId のままだと
# 「本当に新しいデータ」でも既知バッチとしてスキップされてしまう。
# 意図的にゼロから作り直すときは別の txnAppId にする。
dbutils.fs.rm(ck_idem, True)
run_idem("orders-ingest-v2")
print("別 txnAppId:", spark.table(f"{CATALOG}.{SCHEMA}.orders_idem").count())          # 2000
python
# --- ④ foreachBatch + MERGE(Delta の update 相当)---
def upsert(micro_df, batch_id):
    if micro_df.isEmpty():
        return
    micro_df.createOrReplaceTempView("updates")
    micro_df.sparkSession.sql(f"""
      MERGE INTO {CATALOG}.{SCHEMA}.orders_agg t
      USING updates s ON s.id = t.id
      WHEN MATCHED THEN UPDATE SET *
      WHEN NOT MATCHED THEN INSERT *
    """)

spark.sql(f"CREATE TABLE IF NOT EXISTS {CATALOG}.{SCHEMA}.orders_agg (id BIGINT, amount BIGINT)")
q = (spark.readStream.table(f"{CATALOG}.{SCHEMA}.raw_orders")
     .writeStream.foreachBatch(upsert)
     .option("checkpointLocation", f"{VOL}/_ckpt/merge")
     .trigger(availableNow=True).start())
q.awaitTermination()
python
# --- ⑤ Delta シンクの出力モード制約を確認する ---
from pyspark.sql.functions import count as _count
agg = spark.readStream.table(f"{CATALOG}.{SCHEMA}.raw_orders").groupBy("amount").agg(_count("*").alias("n"))
try:
    (agg.writeStream.outputMode("update")           # ← Delta は update 非サポート
        .option("checkpointLocation", f"{VOL}/_ckpt/upd")
        .trigger(availableNow=True)
        .toTable(f"{CATALOG}.{SCHEMA}.agg_update")).awaitTermination()
except Exception as e:
    print("Delta は update 非サポート:", str(e)[:200])

確認ポイント

  • foreachBatch は at-least-once のみ保証。exactly-once 相当は txnAppIdtxnVersionbatchId にバインド)
  • チェックポイントを作り直すときは別の txnAppId(新チェックポイントは batchId 0 から始まるため)
  • 空 DataFrame は渡されうるOPTIMIZE で処理対象ファイルがない等)ので必ず処理する
  • foreachBatch 内で複数シンクに書くとシリアル化されてレイテンシが増える → シンクごとに別ライター
  • foreachBatch連続処理モードでは動かない(そちらは foreach
  • Delta シンクは append / complete のみ、update は非サポートforeachBatchMERGE で代替。Kafka は全モード対応

課題 13|チェックポイントを壊す/変えると何が起きるか

目的: 「クエリごとに別チェックポイント」「削除=新規開始」「変えてよい変更/ダメな変更」を確認する。 対応教材: 同 3-3 / 確認問題: Q6

python
ck = f"{VOL}/_ckpt/cp_demo"

def start(df, mode=None, name="cp_demo"):
    w = df.writeStream.queryName(name).option("checkpointLocation", ck).trigger(availableNow=True)
    if mode:
        w = w.outputMode(mode)
    q = w.toTable(f"{CATALOG}.{SCHEMA}.cp_target")
    q.awaitTermination()
    return q

src = spark.readStream.table(f"{CATALOG}.{SCHEMA}.raw_orders")

# 1) 正常実行 → チェックポイントの中身を見る
start(src)
for d in ("offsets", "commits", "metadata", "state"):
    try:
        print(d, dbutils.fs.ls(f"{ck}/{d}")[:3])
    except Exception as e:
        print(d, "(なし)")
print(dbutils.fs.head(f"{ck}/metadata"))     # 一意のクエリ ID
python
# 2) 変えてよい変更:フィルターの追加・レート制限・トリガー間隔
from pyspark.sql.functions import col
start(src.filter(col("amount") > 500))       # フィルター追加は OK(同じチェックポイントで再開)

# 3) 変えてはいけない変更:ステートフル操作の追加(状態スキーマが変わる)
from pyspark.sql.functions import count as _count
try:
    start(src.groupBy("amount").agg(_count("*")), mode="complete")
except Exception as e:
    print("状態スキーマ非互換:", str(e)[:200])

# 4) チェックポイントを削除 → 新規開始(既存データを全部読み直す)
dbutils.fs.rm(ck, True)
spark.sql(f"TRUNCATE TABLE {CATALOG}.{SCHEMA}.cp_target")
start(src)
print("再取り込み:", spark.table(f"{CATALOG}.{SCHEMA}.cp_target").count())

# 5) 2 つのクエリで同じチェックポイントを共有してはいけないことを確認
python
# 6) 後片付け(クォータ節約)
for q in spark.streams.active:
    print("stopping", q.name); q.stop()

確認ポイント

  • チェックポイントの 4 要素 = offsets / commits / state / metadata(クエリ ID)
  • クエリごとに別の場所が必須削除/場所変更=新規開始
  • 安全な変更: フィルター追加/削除、レート制限、トリガー間隔mapGroupsWithState の UDF ロジック
  • 新チェックポイントが必要: 入力ソースの数・種類、購読 Kafka トピック/Auto Loader パス、ステートフル操作の種類・状態スキーマ出力シンクの種類
  • Free Edition の制約: トリガーは AvailableNow / Once のみ。ProcessingTime への変更で「トリガー変更は安全」を試すことはできない

課題 14|リネージ・タグ・BROWSE とガバナンス運用

目的: 列レベルリネージとタグ制約を実機で確認する。 対応教材: 08-governance.md確認問題: 08-governance.md

sql
-- 1) リネージを作る(列レベルまで自動キャプチャされる)
CREATE OR REPLACE TABLE lin_bronze AS SELECT * FROM samples.nyctaxi.trips LIMIT 5000;
CREATE OR REPLACE TABLE lin_silver AS
  SELECT tpep_pickup_datetime AS pickup, fare_amount AS fare, trip_distance AS dist FROM lin_bronze;
CREATE OR REPLACE TABLE lin_gold AS
  SELECT to_date(pickup) AS d, sum(fare) AS revenue FROM lin_silver GROUP BY 1;

-- Catalog Explorer で lin_gold →[系列]→[系列グラフの表示]→ revenue 列をクリックして
-- fare → revenue の列レベル系列を確認する

-- 2) システムテーブルからリネージを引く(Free Edition ではアクセス不可の可能性あり)
SELECT source_table_full_name, target_table_full_name, entity_type, event_time
FROM system.access.table_lineage
WHERE target_table_full_name = '<catalog>.de_sandbox.lin_gold'
  AND event_time >= current_date() - INTERVAL 7 DAYS
ORDER BY event_time DESC;

SELECT source_table_full_name, source_column_name, target_table_full_name, target_column_name
FROM system.access.column_lineage
WHERE target_table_full_name = '<catalog>.de_sandbox.lin_gold';
sql
-- 3) タグの制約を実際に確認する
ALTER TABLE lin_silver SET TAGS ('domain' = 'nyctaxi', 'pii' = 'false');
SELECT * FROM <catalog>.information_schema.table_tags WHERE table_name = 'lin_silver';

-- 大文字小文字を区別する(Sales と sales は別タグ)
ALTER TABLE lin_silver SET TAGS ('Domain' = 'X');
SELECT * FROM <catalog>.information_schema.table_tags WHERE table_name = 'lin_silver';

-- 列タグ(1 コマンド 1 列。複数列同時は不可)
ALTER TABLE lin_silver ALTER COLUMN fare SET TAGS ('pii' = 'false');
SELECT * FROM <catalog>.information_schema.column_tags WHERE table_name = 'lin_silver';

-- キーに使えない文字(. , - = / :)を入れるとエラー
ALTER TABLE lin_silver SET TAGS ('cost-center' = 'hr');   -- ← エラー

-- 4) BROWSE とデータ検出
GRANT BROWSE ON CATALOG <catalog> TO `account users`;
SHOW GRANTS ON CATALOG <catalog>;

確認ポイント

  • リネージは列レベルまで自動キャプチャされ、メタストア配下の全ワークスペースで集約される。表示には最低 BROWSE
  • リネージの制限: 名前変更で保持されない / パス参照では列系列不可 / UDF がマッピングを隠す / RDD・グローバル一時ビュー不可 / 2024-09-01 以前は不可 / システムテーブルはローリング 1 年
  • タグ: 1 オブジェクト最大 50、キー/値 256 文字キーに . , - = / : 不可大文字小文字を区別1 コマンドで複数列は不可平文で保存されるので機密情報を入れない
  • ABAC のときだけ上位からタグが暗黙継承される(列レベルは除く)
  • Free Edition の制約: システムテーブル(system.access.*)はアカウント管理者+メタストア管理者前提のためアクセスできない可能性が高い。エラーになったら座学に切り替える。管理タグ(ASSIGN 特権)もアカウントレベル機能

課題 15|SCD Type 1/2・AUTO CDC・CDF を作り比べる

目的: 同じ CDC 入力から Type 1 と Type 2 を作り、順序外イベントの結果を自分の目で確認する。教材の公式例をそのまま再現。 対応教材: 09-data-modeling.md確認問題: 09-data-modeling.md

15-1. CDC ソースを作る(順序外を含む)

sql
CREATE OR REPLACE TABLE users_cdf (
  userId INT, name STRING, city STRING, operation STRING, sequenceNum INT
);
INSERT INTO users_cdf VALUES
  (124,'Raul','Oaxaca','INSERT',1),
  (123,'Isabel','Monterrey','INSERT',1),
  (125,'Mercedes','Tijuana','INSERT',2),
  (126,'Lily','Cancun','INSERT',2),
  (123,NULL,NULL,'DELETE',6),
  (125,'Mercedes','Guadalajara','UPDATE',6),
  (125,'Mercedes','Mexicali','UPDATE',5),      -- 順序外
  (123,'Isabel','Chihuahua','UPDATE',5);       -- 順序外

15-2. パイプラインで SCD Type 1 と Type 2 を作る

sql
-- Type 1:最新のみ
CREATE OR REFRESH STREAMING TABLE users_current;

CREATE FLOW apply_cdc_t1 AS AUTO CDC INTO users_current
FROM stream(<catalog>.de_sandbox.users_cdf)
KEYS (userId)
APPLY AS DELETE WHEN operation = "DELETE"
SEQUENCE BY sequenceNum
COLUMNS * EXCEPT (operation, sequenceNum)
STORED AS SCD TYPE 1;

-- Type 2:全履歴
CREATE OR REFRESH STREAMING TABLE users_history;

CREATE FLOW apply_cdc_t2 AS AUTO CDC INTO users_history
FROM stream(<catalog>.de_sandbox.users_cdf)
KEYS (userId)
APPLY AS DELETE WHEN operation = "DELETE"
SEQUENCE BY sequenceNum
COLUMNS * EXCEPT (operation, sequenceNum)
STORED AS SCD TYPE 2;

パイプラインを更新したら結果を検算します。

sql
SELECT * FROM users_current ORDER BY userId;
-- 期待:124 Oaxaca / 125 Guadalajara / 126 Cancun(123 は削除、seq=5 の Mexicali は seq=6 に負ける)

SELECT * FROM users_history ORDER BY userId, __START_AT;
-- 期待:123 Monterrey(1→5), 123 Chihuahua(5→6), 124 Oaxaca(1→null),
--       125 Tijuana(2→5), 125 Mexicali(5→6), 125 Guadalajara(6→null), 126 Cancun(2→null)

15-3. TRACK HISTORY ON ... EXCEPT で追跡列を絞る

sql
CREATE OR REFRESH STREAMING TABLE users_history_no_city;

CREATE FLOW apply_cdc_t2b AS AUTO CDC INTO users_history_no_city
FROM stream(<catalog>.de_sandbox.users_cdf)
KEYS (userId)
APPLY AS DELETE WHEN operation = "DELETE"
SEQUENCE BY sequenceNum
COLUMNS * EXCEPT (operation, sequenceNum)
STORED AS SCD TYPE 2
TRACK HISTORY ON * EXCEPT (city);
-- 期待:city の変更では新履歴を作らないので 123 は 1 行、125 は 1 行に集約される

15-4. CDF(変更データフィード)

sql
-- レガシー CDF を有効化して行レベル変更を読む
ALTER TABLE <catalog>.de_sandbox.dim_customer
  SET TBLPROPERTIES (delta.enableChangeDataFeed = true);

UPDATE <catalog>.de_sandbox.dim_customer SET city = 'Kobe' WHERE id = 1;
DELETE FROM <catalog>.de_sandbox.dim_customer WHERE id = 3;

-- _change_type に insert / update_preimage / update_postimage / delete が出る
SELECT _change_type, _commit_version, _commit_timestamp, *
FROM table_changes('<catalog>.de_sandbox.dim_customer', 0)
ORDER BY _commit_version;

-- 開始バージョンを省略するとエラーになることを確認(バッチ読み取りでは必須)
SELECT * FROM table_changes('<catalog>.de_sandbox.dim_customer');
python
# ストリーミング読み取り:初回は最新スナップショットを INSERT として返す
(spark.readStream
  .option("readChangeFeed", "true")
  .table(f"{CATALOG}.{SCHEMA}.dim_customer")
  .writeStream
  .option("checkpointLocation", f"{VOL}/_ckpt/cdf")
  .trigger(availableNow=True)
  .toTable(f"{CATALOG}.{SCHEMA}.dim_customer_changes")).awaitTermination()

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

# 既存状態を INSERT として再処理したくないときは startingVersion を指定する

確認ポイント

  • Type 1 は上書き・最新のみ / Type 2 は __START_AT / __END_AT で全履歴、アクティブ行は __END_AT = NULLシーケンス値がそのまま __START_AT / __END_AT に伝播する
  • 順序外イベントはシーケンス値で決定論的に処理されるseq=5seq=6 に負ける)
  • TRACK HISTORY ON ... EXCEPT (...) の対象外列の変更は新履歴を作らず上書き
  • CDF=変更の記録/読み取り、AUTO CDC=変更適用による SCD 具体化。役割を混同しない
  • CDF のメタデータ列は _change_type / _commit_version / _commit_timestampバッチ読み取りは開始バージョン必須
  • CDF は恒久記録ではない(保持期間内のみ、非加法スキーマ変更を跨げない)
  • Free Edition の制約: AUTO CDCPro / Advanced エディション(またはサーバーレス)のパイプラインが必要。サーバーレスパイプラインなので動くはずですが、アクティブ 1 本の制約があるため Type 1 / Type 2 を 1 つのパイプラインにまとめて作るのが効率的。AUTO CDC FROM SNAPSHOT は Python のみなので Python ファイルで別途試す

ドメイン⑩ データ共有・フェデレーションについて

Free Edition では実質練習できません。

  • Delta Sharing: プロバイダー側の受信者作成・トークン管理は有料エディション向けとされ、Free Edition では制限されます。受信者側も相手のプロバイダーが必要です
  • Lakehouse フェデレーション: 接続先の外部 DB(PostgreSQL / MySQL / Snowflake 等)とアウトバウンド通信が必要で、Free Edition はアウトバウンドが制限されます

配点は 5%(59 問中およそ 3 問)なので、次の 3 点セットで確実に取りに行きます。

  1. 教材 10-sharing-federation.md を通読
  2. 確認問題 10-sharing-federation.md(10 問)を満点まで回す
  3. 構文だけは書けるようにする(CREATE SHARE / ALTER SHARE ADD TABLE ... WITH HISTORY / GRANT SELECT ON SHARE ... TO RECIPIENT / CURRENT_METASTORE() / CURRENT_RECIPIENT() / CREATE CONNECTIONCREATE FOREIGN CATALOGGRANT

なお SELECT CURRENT_METASTORE(); は Free Edition でも実行できるので、共有識別子の形式(<cloud>:<region>:<uuid>)だけは実機で確認しておきましょう。


仕上げチェック

Free Edition で確かめられなかった論点の最終確認リスト(すべて 00-free-edition.md 6 章の暗記表と対応)

  • [ ] Spark UI の診断順序(タイムライン → 最長ステージ → スピル → スキュー → I/O → その他)とスキュー 50% ルール
  • [ ] Executor 削除の 3 理由(自動スケール / スポット損失 / OOM)と調べ方(イベントログ →[Executors]タブ)
  • [ ] クラスターサイズ設計・インスタンスファミリの対応・スポット(ドライバーはオンデマンド)・プール(アイドル中は DBU 課金なし)
  • [ ] AQE の 4 機能と既定値、結合順序は変えない
  • [ ] トリガー種別の使い分け(ProcessingTime / AvailableNow / realTime、Continuous は非サポート)
  • [ ] Scala UDF の位置(組み込み > Scala > pandas > 行ごと Python)
  • [ ] Delta Sharing(オープン vs D2D、認証、WITH HISTORY)と Lakehouse フェデレーション(2 タイプ、3 ステップ、読み取り専用)
  • [ ] システムテーブルのパスと権限(USE CATALOG + USE SCHEMA + SELECT)、リアルタイム監視不可
  • [ ] ABAC ポリシーの句(WHEN has_tag_value / MATCH COLUMNS 最大 3 / TO EXCEPT / USING COLUMNS

最後に片付け

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