Skip to content

Professional 学習教材 ③ データ変換・クレンジング・品質(配点 10%)

本教材は Databricks 認定データエンジニア Professional(上位資格) 試験の「データ変換・クレンジング・品質」ドメインを、公式ドキュメント(日本語版)を実読したうえで、本 md 単体で自習が完結するようにまとめたものです。本番運用(プロダクション ETL)観点でのデータ品質保証・クレンジング・変換を深く扱います。

参照した公式ページ(すべて WebFetch で実読、取得日 2026-07-25):

※ 本文中の API 表記について: 最新ドキュメントでは Lakeflow パイプライン(旧 Delta Live Tables / DLT)の Python デコレーターが from pyspark import pipelines as dp を用いた @dp.table / @dp.expect... 形式で記載されています。旧来の import dlt / @dlt.table / @dlt.expect... と読み替え可能です(試験ではどちらの表記も出題され得ます)。


1. このドメインの概要

このドメインは、レイクハウス上で 信頼できる(trustworthy)データ を作るための、変換(transform)・クレンジング(cleansing)・品質保証(data quality)の実装と運用を問います。Professional では単なる構文暗記ではなく、次の「本番運用の判断」ができることが求められます。

  • 変換のセマンティクス選択: バッチ処理とストリーミング処理の違いを理解し、メダリオン(bronze / silver / gold)の各層に適切な処理方式を選べる。
  • データ品質の作り込み: Lakeflow パイプラインの Expectations(期待値) を使い、無効レコードを「警告(warn)/削除(drop)/失敗(fail)」のどの方針で扱うかを設計できる。無効データの 隔離(quarantine)検証テーブル による制御も実装できる。
  • クレンジング: 欠損値・重複・型不一致・スキーマ変化に対処し、MERGE INTOアップサート(upsert)/重複排除(deduplication)/条件分岐 を実装できる。冪等性(idempotency) を保った書き込みができる。
  • 品質の監視: データ品質の監視(旧 Lakehouse Monitoring) による データプロファイリング異常検知(鮮度・完全性)で、時間経過にともなう品質変化・ドリフトを継続監視できる。

キーワードの相互関係(大枠):

[変換 transform] ── バッチ / ストリーミング / 宣言的(Lakeflow) / 命令的(Spark)

        ├─ [クレンジング] 欠損値・重複・型・スキーマ強制/進化
        │        └─ MERGE INTO(upsert / dedup / SCD / 条件分岐)

        ├─ [インライン品質保証] Expectations(warn/drop/fail, quarantine, 検証テーブル)

        └─ [事後の品質監視] データ品質の監視(プロファイリング / 異常検知=鮮度・完全性)

2. 重要用語集

用語(日本語)English説明
期待値ExpectationsLakeflow パイプラインで、マテリアライズドビュー/ストリーミングテーブル/ビューの作成文に付ける任意の句。通過する各レコードに SQL ブール制約でデータ品質チェックを適用する。1 データセットに複数指定でき、複数データセットで再利用可能。
期待値名Expectation name各期待値に必須の識別子。追跡・監視に使う。特定データセット内で一意である必要がある(例: valid_customer_age)。
制約句Constraint / EXPECT clause各レコードで true/false に評価される SQL 条件式。false のレコードで期待値がトリガーされる。カスタム Python 関数・外部サービス呼び出し・他テーブル参照のサブクエリは使用不可。
保持(警告)expect(warn, 既定値)既定動作。違反レコードもターゲットに書き込み、合否件数のメトリックを収集。SQL は EXPECT、Python は dp.expect
削除expect_or_drop(drop)違反レコードをターゲット書き込み前に削除。削除件数をログに記録。SQL は EXPECT ... ON VIOLATION DROP ROW、Python は dp.expect_or_drop
失敗expect_or_fail(fail)違反レコードがあると更新を失敗させ即停止。テーブル更新はアトミックにロールバック。SQL は EXPECT ... ON VIOLATION FAIL UPDATE、Python は dp.expect_or_fail
違反時アクション句ON VIOLATIONSQL で違反時の挙動を指定する句。DROP ROW(削除)または FAIL UPDATE(失敗)。省略時は warn。
複数期待値の一括指定expect_all / expect_all_or_drop / expect_all_or_failPython 限定。期待値を辞書(名前→制約)でまとめて指定し、グループに集合アクションを適用。複数データセットで同一ルールを再利用できる。
データ品質ルールData quality rules期待値の制約定義群。Delta テーブルや Python モジュール/辞書に外部化して再利用・監査・保守できる(可搬性)。
隔離Quarantineデータを失わずに無効レコードを別経路で扱うパターン。is_quarantined フラグ列+期待値+一時テーブル/ビューで有効/無効を分離する。
検証テーブルValidation table他テーブル間のプロパティ(行数一致・主キー一意性・欠損なし等)をチェックし expect_or_fail で失敗させる別データセット。オーケストレーションではなく品質担保が目的。
MERGE INTOMERGE INTOソースを基にターゲット Delta テーブルへ更新・挿入・削除をまとめて適用する文。Delta Lake テーブルのみ対応。
アップサートUpsert既存キーは更新、なければ挿入する操作。WHEN MATCHED THEN UPDATEWHEN NOT MATCHED THEN INSERT で実現。
重複排除Deduplication重複レコードを除去する処理。MERGE の insert-only(WHEN NOT MATCHED THEN INSERT)で既存キーの再挿入を防ぐ。
スキーマ強制Schema enforcement書き込むデータがテーブルスキーマに一致することを強制し、不一致を拒否する仕組み(Delta の既定)。
スキーマ進化Schema evolution新しい列などスキーマ変化を自動的にテーブルへ反映する仕組み。MERGE WITH SCHEMA EVOLUTION、Auto Loader の cloudFiles.schemaLocation 等。
冪等性Idempotency同じ処理を何度実行しても結果が変わらない性質。MERGE の insert-only 重複排除やチェックポイントにより、再実行・重複ログでも二重書き込みを防ぐ。
データプロファイリングData profilingテーブルの概要統計(null 率・分布・パーセンタイル等)を時系列でキャプチャし、品質・一貫性の変化を追跡する(旧 Lakehouse Monitoring)。
データ品質の監視Data quality monitoringUnity Catalog 上の資産の品質を確保する機能群。異常検知データプロファイル を含む。サーバーレスで実行。
異常検知Anomaly detectionスキーマ内全テーブルをスケーラブルに監視し、鮮度(freshness)完全性(completeness) を自動評価。インテリジェントスキャンで重要テーブルを優先。
鮮度Freshnessテーブルが最近更新されたか。コミット履歴からテーブルごとに次コミット時刻を予測し、異常に遅れると「古い(stale)」とマーク。
完全性Completeness過去 24 時間の書き込み行数。履歴から予測範囲を出し、下限を下回ると「不完全(incomplete)」とマーク。
ドリフトDrift(誤差/偏移)現在データと既知のベースライン間、または連続する時間枠間の統計的なずれ。プロファイルで検出できる。
欠損値処理Missing value handlingnull/欠損への対処。IS NOT NULL 期待値、dropMERGEDEFAULT、フィルタ等で扱う。
バッチ処理Batch processingソースの利用可能な全データを毎回再処理。ロジックが単純で結果は常に正確だが非効率・低速。
ストリーミング処理Streaming processing前回以降の新データのみを追跡・処理。効率的・低遅延だが、順序ずれ・遅延データ・ステートフル処理で複雑化。
宣言的変換Declarative transformationLakeflow パイプラインのように「何を作るか」を宣言し、増分処理・依存関係・品質適用を基盤に委ねる方式。
自動ローダーAuto Loaderクラウドオブジェクトストレージに到着した新ファイルを自動検出・処理する増分インジェスト機能(format("cloudFiles"))。
メダリオンアーキテクチャMedallion architecturebronze(取り込み)→ silver(変換・クレンジング)→ gold(集約)でデータを段階的に洗練する設計。

3. 詳細解説

3-1. 変換の全体像(バッチ / ストリーミング、宣言的変換)

バッチとストリーミングのセマンティクス

Databricks では、バッチとストリーミングは「インジェスト・変換・リアルタイム処理」に使う 2 つの 処理セマンティクスです。Lakeflow パイプライン(Apache Spark + Structured Streaming)は両者を統合したアーキテクチャを持ちます。

  • バッチセマンティクス: エンジンは処理済みデータを追跡しない。ソースで現在利用可能な全データを毎回処理する。実務では再処理を抑えるため日付・地域などで論理パーティションすることが多い。例: 1 時間ごとに前の時間の全データを再処理し、以前の結果を上書きして最新化。
  • ストリーミングセマンティクス: エンジンは処理済みを追跡し、後続実行では新データのみを処理する。完全な結果を得るには過去の結果に新規結果を追加していく必要がある。

トレードオフ(重要な比較):

セマンティクス利点短所主なデータエンジニアリング製品
バッチ処理ロジックが単純。結果は常に正確(全データ反映)。非効率(再処理)。低速で、時間〜分単位は可能だが秒・ミリ秒は不可。Lakeflow の具体化ビュー/具体化ビューフロー。Databricks Runtime の Apache Spark(spark.read.load() / spark.write.save())。
ストリーミング効率的(新データのみ)。高速で、時間〜分〜秒〜ミリ秒まで対応。結合・集計・重複除去などステートフル処理で複雑化。順序ずれ・遅延データで結果が常に正確とは限らない。Lakeflow Connect、Lakeflow パイプライン(追加/変更フロー、ストリーミングテーブル、シンク)、Structured Streaming(spark.readStream.load() / spark.writeStream.start())。

到着遅延データ(late-arriving data)の扱いの違い(試験頻出の考え方):

  • バッチ: 遅延データは次回バッチで既存データと合わせて再処理され、前回結果は上書きで修正される。
  • ストリーミング: 遅延データはその時間帯の既処理データとは別に処理される。前の結果を正しく更新するには、合計・件数などの 状態(state) を保持する必要がある。

ステートレス/ステートフル:

  • ステートレス(フィルタ、単純追加)は順序ずれ・遅延に強く、複雑にならない。
  • ステートフル(結合 join・集計 aggregation・重複除去 dedup)はウォーターマークなどの状態管理が必要で複雑になる。

メダリオン各層での推奨処理

ワークロード特性推奨
Bronze(銅)取り込み。多くはステートレスな増分追加。データサイズ大。ストリーミング処理(ステートフルの複雑さに触れずに増分の利点を得られる)。
Silver(銀)変換。フィルタ等のステートレス+結合・集計・重複除去などステートフルの両方。バッチ処理(具体化ビューは増分更新)。効率・遅延が精度より重要ならストリーミングも選択肢(複雑さに注意)。
Gold(金)最終段の集約。結合・集計などステートフル。データサイズ小。バッチ処理(具体化ビューは増分更新)。

宣言的変換 vs 命令的変換

  • 宣言的(declarative): Lakeflow パイプラインで、テーブル/ビューの定義と品質(Expectations)を宣言。増分処理・依存解決・データ品質適用は基盤が担う。運用 ETL の構築・デプロイ・保守の複雑さが減る。
  • 命令的(imperative): Databricks Runtime 上の Apache Spark / Structured Streaming を直接記述(spark.read / readStream、DataFrame 変換、writeStream 等)。ETL クイックスタートはこの方式。

3-2. Expectations によるデータ品質保証(warn / drop / fail、複数ルール)

Expectations の 3 要素

期待値は次の 3 つで構成されます。

  1. 期待値名: 一意の識別子。検証内容が伝わる名前を付ける。データセット内で一意(複数データセット間では再利用可)。
  2. 評価する制約(EXPECT 句): 各レコードで true/false になる SQL 条件。制約に使えないもの=カスタム Python 関数、外部サービス呼び出し、他テーブルを参照するサブクエリ。
  3. 違反時アクション: warn / drop / fail のいずれか。

3 つのアクション(最重要の比較表)

アクションSQL 構文Python 構文結果
warn(既定)EXPECTdp.expect無効レコードもターゲットに書き込まれる。合否件数のメトリックを収集。
drop(削除)EXPECT ... ON VIOLATION DROP ROWdp.expect_or_drop無効レコードは書き込み前に削除。削除件数をメトリックに記録。
fail(失敗)EXPECT ... ON VIOLATION FAIL UPDATEdp.expect_or_fail無効レコードがあると更新失敗。トランザクションはアトミックにロールバック。再処理には手動介入が必要。
  • 既定は warn(保持)。合否件数を測りたいだけならこれを使う。違反レコードは有効レコードと共にターゲットに追加される。
  • drop は違反レコードを以降の処理に流さない。
  • fail は違反が許容できない場合に即停止。メトリックは記録されない(更新が失敗するため)。

fail 時の挙動(パイプライン実行モード依存・頻出)

  • トリガーされたパイプライン: 1 フローが失敗しても、他の並列フローは失敗しない。失敗したフローの更新のみロールバックされる。
  • 連続(continuous)パイプライン: 失敗した期待値がフローを停止し、すべての依存フローも停止。停止理由のメッセージを出力。
  • fail 用に構成された期待値は、違反を検出・報告するために Spark クエリプランを変更し、違反した入力レコードを特定できるようにする。専用エラーの例:
[EXPECTATION_VIOLATION.VERBOSITY_ALL] Flow 'sensor-pipeline' failed to meet the expectation. Violated expectations: 'temperature_in_valid_range'. Input data: '{"id":"TEMP_001","temperature":-500,"timestamp_ms":"1710498600"}'. Output record: '{"sensor_id":"TEMP_001","temperature":-500,"change_time":"2024-03-15 10:30:00"}'. Missing input data: false
  • 検証失敗時にダウンストリームの制御をより細かくしたい場合は、検証とダウンストリーム作業を 別々のパイプラインに分割 し、ジョブ(タスク依存)で調整する。

期待値メトリックの追跡

  • warn / drop のメトリックはパイプライン UI の [データ品質]タブ で確認できる。fail はメトリックが記録されない。
  • 手順: サイドバーの[ジョブ & パイプライン]→ パイプライン名 → 期待値が定義されたデータセット → 右サイドバーの[データ品質]タブ。
  • Databricks SQL で作成したスタンドアロンパイプラインでは[データ品質]タブは使えず、イベントログにクエリしてメトリックを取得する。

複数の期待値の管理

  • SQL・Python とも 1 データセットに複数の期待値を持てる。ただし グループ化して集合アクションを指定できるのは Python のみexpect_all / expect_all_or_drop / expect_all_or_fail、引数は「名前→制約」の辞書)。同じルールセットを複数データセットで再利用できる。

期待値の制限事項

  • 期待値をサポートするのは ストリーミングテーブル・具体化ビュー・一時ビュー のみ(データ品質メトリックもこれらだけ)。
  • メトリックが得られないケース: 期待値未定義、期待値非対応の演算子/フロー種別(例: シンク)、当該実行で対象テーブルへの更新がない、pipelines.metrics.flowTimeReporter.enabled 等の設定不足、ビューは実行時のみ計算されるため取得できない/複数セットになる場合がある。
  • 期待値は AUTO CDC FROM SNAPSHOT では非対応。
  • Databricks SQL 作成のスタンドアロンパイプラインでも、CREATE STREAMING TABLE / CREATE MATERIALIZED VIEWCONSTRAINT ... EXPECT (...) で期待値を定義できる。

3-3. クレンジング(欠損値・重複・型不一致・スキーマ強制/進化)

クレンジングは silver 層の中心作業です。Databricks では次の手段を組み合わせます。

欠損値(null)処理

  • 期待値で IS NOT NULL を検証(例: current_page_id IS NOT NULL AND current_page_title IS NOT NULL)。
    • 単に測る → warn、除外 → expect_or_drop、致命的 → expect_or_fail
  • MERGEINSERT (...) VALUES (...) でターゲット列を省略すると、列の 既定値 が、無ければ NULL が挿入される。DEFAULT を明示指定して既定値に更新/挿入することもできる。
  • INSERT * EXCEPT (col) / UPDATE SET * EXCEPT (col) で除外した列は null に設定される(スキーマ進化有効時はソース列参照&進化対象外に変わる)。

重複(duplication)排除

  • 詳細は 3-4 の MERGE を参照。核心は insert-only マージWHEN NOT MATCHED THEN INSERT)で既存キーの再挿入を防ぐこと。
  • 重要な注意: MERGE はターゲット既存データとの重複は防ぐが、新規データセット内部の重複は挿入されてしまう。したがって マージ前に新データ側を重複排除しておく必要がある。
  • ストリーミングの重複除去はウォーターマーク(watermarks#drop-duplicates-within-watermark)等でステートフルに扱う。

型不一致(type mismatch)

  • 期待値で SQL 関数を使い検証(例: year(transaction_date) >= 2020、日付キャスト to_date(updateTime,'M/d/yyyy h:m:s a') > '2010-01-01')。
  • CASE 式で型・状態に応じた条件分岐検証が可能。

スキーマ強制(enforcement)とスキーマ進化(evolution)

  • 強制: Delta は既定でスキーマを強制し、一致しない書き込みを拒否する(データ整合性を守る)。
  • 進化(MERGE): MERGE WITH SCHEMA EVOLUTION(Databricks Runtime 15.2 以上)で、ターゲットのスキーマをソースに合わせて自動更新。UPDATE SET * / INSERT * は「ソース列=ターゲット列」を前提とし、異なると分析エラー。
  • 進化(インジェスト): Auto Loader は cloudFiles.schemaLocation(=チェックポイント位置)でスキーマを推論・追跡し、進化を扱う。
  • 進化パターン(複数バージョン混在): 旧 V1 と新 V2 を統合する際、Python は unionByName(new_data, allowMissingColumns=True)、SQL は SELECT *, NULL as col3 FROM legacy_sourceUNION で結合し、期待値で後方互換を担保しつつ品質を適用する(3-2 の複数ルールと併用)。

3-4. MERGE INTO による変換(アップサート・重複排除・条件分岐)

MERGE はソースを基にターゲット Delta テーブルへ 更新・挿入・削除 をまとめて適用します(Delta テーブル専用、外部テーブル不可)。SCD タイプ 2、CDC のアップサート、重複除去など複雑な操作に使えます。

完全構文

[ 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 [AND cond]ソース行がターゲット行に一致(ON +任意条件)DELETE / UPDATE SET * / UPDATE SET col=...更新・削除
WHEN NOT MATCHED [BY TARGET] [AND cond]ソース行がどのターゲット行にも一致しないINSERT * / INSERT (...) VALUES (...) のみ挿入
WHEN NOT MATCHED BY SOURCE [AND cond]ターゲット行がどのソース行にも一致しない(DBR 12.2 LTS 以上)DELETE / UPDATE SET col=...削除・更新

句のルール(頻出):

  • 各種類の句は複数指定でき、指定順に評価される。
  • 同じ種類の句が複数ある場合、最後の句を除くすべてに条件(AND …)が必須。省略すると NON_LAST_..._CLAUSE_OMIT_CONDITION エラー。
  • WHEN NOT MATCHED BY SOURCE の条件・アクションは ターゲット列のみ参照可(ソース行が存在しないため)。リテラルやターゲット列演算(例: SET target.deleted_count = target.deleted_count + 1)で指定する。
  • WHEN NOT MATCHED BY SOURCE を条件なしで使うと、多数のターゲット行の更新/削除・テーブル全書き換えにつながりコスト増。必ず条件で範囲を絞るのが推奨。
  • 複数一致(multiple matches)は失敗: ソースの複数行がターゲットの同一行に一致すると操作は失敗(どのソース行で更新すべきか曖昧なため)。ソースを前処理して一致を一意化する(例: キーごとに最新変更のみ残す=CDC 前処理)。DBR 16.0 以降は ONWHEN MATCHED 条件で重複一致を判定、15.4 LTS 以下は ON 条件のみで判定。ただし無条件の DELETE は複数一致でも曖昧でないため許可される。

UPDATE SET * / INSERT * の意味

  • UPDATE SET * = 全列を col = source.col で更新。ソース列とターゲット列が同一である前提(違うと分析エラー)。
  • INSERT * = 全列を INSERT (cols) VALUES (source.cols) で挿入。同前提。
  • EXCEPT (col) で特定列を除外(除外列は null、スキーマ進化時は挙動が変わる)。
  • DEFAULT を式に指定して既定値へ更新/挿入できる。

パターン別コード例(本文の 4「構文・コード例」に集約)

アップサート、insert-only 重複排除、WHEN NOT MATCHED BY SOURCE による同期、SCD/CDC の考え方は 4 章にまとめています。

冪等性(idempotency)

  • insert-only マージWHEN NOT MATCHED THEN INSERT *)は同じキーの再挿入を防ぐため、再実行や重複ログでも二重書き込みされず冪等になる。ログ追記型 ETL の重複対策の定番。
  • 数日だけ重複が発生し得る場合は、テーブルを日付でパーティションし、ONWHEN NOT MATCHED に日付範囲条件を付けて一致範囲を限定するとクエリを最適化できる(テーブル全体でなく直近のみ探索)。
  • insert-only マージは新データのみ追加するため、foreachBatch と組み合わせた 構造化ストリーミングの継続的重複排除 にも使え、別ストリームから重複排除済みデータを継続読み取りできる。

3-5. データ品質監視(プロファイリング・異常検知)

データ品質の監視(Data quality monitoring、旧 Lakehouse Monitoring) は Unity Catalog 上の全データ資産の品質を確保する機能で、監視対象テーブルを変更せず、設定ジョブにオーバーヘッドを追加しません。サーバーレスコンピューティング で実行され、DATA_QUALITY_MONITORING 課金製品として課金されます。次の 2 機能を含みます。

(A) 異常検知(Anomaly detection)

  • ワンクリックでスケーラブルに、スキーマ内の全テーブルを監視。インテリジェントスキャンで重要テーブルを優先、影響の小さいテーブルはスキップしてコストを抑える(カタログ/スキーマレベル)。
  • 履歴データパターンを分析し、各テーブルの 鮮度完全性 を自動評価する。
    • 鮮度(freshness): テーブルが最近更新されたか。コミット履歴を分析しテーブルごとのモデルで次コミット時刻を予測。異常に遅れると 古い(stale) とマーク。
    • 完全性(completeness): 過去 24 時間に書き込まれる予定の行数。履歴行数から予測範囲を出し、直近 24 時間のコミット行数が下限を下回ると 不完全(incomplete) とマーク。
  • アカウント管理者/メタストア管理者は ガバナンスハブ の[データ]ページで、正常テーブルの割合などの品質ヘルスを確認できる。

(B) データプロファイル(Data profiling/旧 Lakehouse Monitoring)

  • テーブル内データの 概要統計 を提供。テーブルのデータ分布や、対応するモデル性能の 履歴メトリック をキャプチャし、監視・変更アラートに使える。
  • 推論テーブル(モデルの入力・予測)を監視すれば、GenAI アプリ・ML モデル・モデルサービングエンドポイントの性能も追跡できる。
  • 答えられる問いの例(テーブル/リスト):
    • データの整合性は時間でどう変わるか(null/ゼロ値の割合は増えているか)。
    • 統計分布は時間でどう変わるか(数値列の 90 パーセンタイル、カテゴリ列の分布と昨日との差)。
    • 現在データと既知ベースライン間、または連続する時間枠間に 誤差(ドリフト) はあるか。
    • データのサブセット/スライスの分布・ドリフトはどうか。
    • ML モデルの入力・予測は時間でどう変わるか。モデルバージョン A と B の性能比較。
  • 観測の 時間粒度 を制御でき、カスタムメトリック を設定できる。

期待値(インライン)との違い(Professional の勘所)

  • Expectations はパイプライン内で通過時に品質を作り込む(不良を warn/drop/fail、隔離)。
  • データ品質の監視 はテーブルに対する事後・継続的な観測(鮮度・完全性・分布・ドリフト)。テーブルを変更しない。
  • 運用では両者を併用: パイプラインで Expectations により不良を制御し、監視でテーブルの経時的な品質劣化・異常を検知する。

4. 構文・コード例

4-1. Expectations(Python / SQL)

単一の期待値(名前+制約)

Python:

python
@dp.table
@dp.expect("valid_customer_age", "age BETWEEN 0 AND 120")
def customers():
  return spark.readStream.table("datasets.samples.raw_customers")

SQL:

sql
CREATE OR REFRESH STREAMING TABLE customers(
  CONSTRAINT valid_customer_age EXPECT (age BETWEEN 0 AND 120)
) AS SELECT * FROM STREAM(datasets.samples.raw_customers);

3 つのアクション

保持(warn, 既定):

python
@dp.expect("valid timestamp", "timestamp > '2012-01-01'")
sql
CONSTRAINT valid_timestamp EXPECT (timestamp > '2012-01-01')

削除(drop):

python
@dp.expect_or_drop("valid_current_page", "current_page_id IS NOT NULL AND current_page_title IS NOT NULL")
sql
CONSTRAINT valid_current_page EXPECT (current_page_id IS NOT NULL and current_page_title IS NOT NULL) ON VIOLATION DROP ROW

失敗(fail):

python
@dp.expect_or_fail("valid_count", "count > 0")
sql
CONSTRAINT valid_count EXPECT (count > 0) ON VIOLATION FAIL UPDATE

複雑な制約(CASE / 複数条件 / 複合ブール) ※ SQL 関数・CASE・AND/OR が使える

Python:

python
# SQL 関数
@dp.expect("valid_date", "year(transaction_date) >= 2020")

# CASE 文
@dp.expect("valid_order_status", """
   CASE
     WHEN type = 'ORDER' THEN status IN ('PENDING', 'COMPLETED', 'CANCELLED')
     WHEN type = 'REFUND' THEN status IN ('PENDING', 'APPROVED', 'REJECTED')
     ELSE false
   END
""")

# 複数制約(デコレーターを重ねる)
@dp.expect("non_negative_price", "price >= 0")
@dp.expect("valid_purchase_date", "date <= current_date()")

# 複合ブール
@dp.expect("valid_order_state", """
   (status = 'ACTIVE' AND balance > 0)
   OR (status = 'PENDING' AND created_date > current_date() - INTERVAL 7 DAYS)
""")

SQL(複数制約はカンマ区切り):

sql
CONSTRAINT non_negative_price EXPECT (price >= 0),
CONSTRAINT valid_purchase_date EXPECT (date <= current_date())

複数期待値の一括(Python 限定・集合アクション)

python
valid_pages = {"valid_count": "count > 0", "valid_current_page": "current_page_id IS NOT NULL AND current_page_title IS NOT NULL"}

@dp.table
@dp.expect_all(valid_pages)            # warn 集合
def raw_data():
  ...

@dp.table
@dp.expect_all_or_drop(valid_pages)    # drop 集合
def prepared_data():
  ...

@dp.table
@dp.expect_all_or_fail(valid_pages)    # fail 集合
def customer_facing_data():
  ...

4-2. データ品質ルールの外部化(可搬性・再利用)

Delta テーブルにルールを保持

sql
CREATE OR REPLACE TABLE rules AS SELECT
  col1 AS name, col2 AS constraint, col3 AS tag
FROM (VALUES
  ("website_not_null","Website IS NOT NULL","validity"),
  ("fresh_data","to_date(updateTime,'M/d/yyyy h:m:s a') > '2010-01-01'","maintained"),
  ("social_media_access","NOT(Facebook IS NULL AND Twitter IS NULL AND Youtube IS NULL)","maintained")
)

タグでルールを読み込み、辞書として適用get_rules()

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

def get_rules(tag):
  """tag に一致するルールの辞書を rules テーブルから読み込む"""
  df = spark.read.table("rules").filter(col("tag") == tag).collect()
  return {row['name']: row['constraint'] for row in df}

@dp.table
@dp.expect_all_or_drop(get_rules('validity'))
def raw_farmers_market():
  return (spark.read.format('csv').option("header", "true")
    .load('/databricks-datasets/data.gov/farmers_markets_geographic_data/data-001/'))

@dp.table
@dp.expect_all_or_drop(get_rules('maintained'))
def organic_farmers_market():
  return spark.read.table("raw_farmers_market").filter(expr("Organic = 'Y'"))
  • ルールは Python モジュール(例 rules_module.py)に置いても同様に扱える。
  • 注意: SQL では、ファイルからの期待値の動的読み込みは非対応。
  • 推奨(可搬性): 期待定義をパイプラインロジックと分離して保持/カスタムタグでグループ化してフィルタ/類似データセット間で一貫適用。

4-3. 無効レコードの隔離(quarantine)

is_quarantined フラグ列+期待値+一時テーブル/ビューで、有効・無効を分離しつつ品質メトリックを追跡する。

Python:

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

rules = {
  "valid_pickup_zip": "(pickup_zip IS NOT NULL)",
  "valid_dropoff_zip": "(dropoff_zip IS NOT NULL)",
}
quarantine_rules = "NOT({0})".format(" AND ".join(rules.values()))

@dp.view
def raw_trips_data():
  return spark.readStream.table("samples.nyctaxi.trips")

@dp.table(temporary=True, partition_cols=["is_quarantined"])
@dp.expect_all(rules)
def trips_data_quarantine():
  return spark.readStream.table("raw_trips_data").withColumn("is_quarantined", expr(quarantine_rules))

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

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

SQL(期待値を 1 つに統合するか個別にするかを選べる):

sql
CREATE OR REFRESH TEMPORARY STREAMING TABLE trips_data_quarantine(
  -- Option 1: 統合してイベントログ上の名前を 1 つにする
  CONSTRAINT quarantined_row EXPECT (pickup_zip IS NOT NULL OR dropoff_zip IS NOT NULL),
  -- Option 2: 個別に保持し、名前別に複数エントリを残す
  CONSTRAINT invalid_pickup_zip EXPECT (pickup_zip IS NOT NULL),
  CONSTRAINT invalid_dropoff_zip EXPECT (dropoff_zip IS NOT NULL)
)
PARTITIONED BY (is_quarantined)
AS SELECT *, NOT ((pickup_zip IS NOT NULL) and (dropoff_zip IS NOT NULL)) as is_quarantined
   FROM STREAM(raw_trips_data);

4-4. 検証テーブル(expect_or_fail で品質をゲート)

行数の一致検証(変換でデータが失われていないか)

python
@dp.materialized_view(name="count_verification", comment="Validates equal row counts between tables")
@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 table_a),
      (SELECT COUNT(*) AS b_count FROM table_b)""")

主キーの一意性検証

python
@dp.materialized_view(name="report_pk_tests", comment="Validates primary key uniqueness")
@dp.expect_or_fail("unique_pk", "num_entries = 1")
def validate_pk_uniqueness():
  return (spark.read.table("report").groupBy("pk").count()
    .withColumnRenamed("count", "num_entries"))

欠損レコード検出(左外部結合でターゲットに存在するか)

python
@dp.materialized_view(name="report_compare_tests", comment="Validates no records are missing after joining")
@dp.expect_or_fail("no_missing_records", "r_key IS NOT NULL")
def validate_report_completeness():
  return (spark.read.table("validation_copy").alias("v")
    .join(spark.read.table("report").alias("r"), on="key", how="left_outer")
    .select("v.*", "r.key as r_key"))

重要な原則: Expectations は オーケストレーションではなく品質担保。検証テーブルを他データセットから読んでも、そのデータセットは検証結果を待たない(ダウンストリームは自動でゲートされない)。検証失敗でダウンストリームを止めたいなら、検証とダウンストリームを 別パイプラインに分け、ジョブのタスク依存で調整する。

4-5. MERGE INTO の各パターン

基本のアップサート(upsert)

sql
MERGE INTO people10m
USING people10mupdates
ON people10m.id = people10mupdates.id
WHEN MATCHED THEN
  UPDATE SET id = people10mupdates.id, firstName = people10mupdates.firstName, /* ... */ salary = people10mupdates.salary
WHEN NOT MATCHED
  THEN INSERT (id, firstName, /* ... */ salary)
  VALUES (people10mupdates.id, people10mupdates.firstName, /* ... */ people10mupdates.salary)

Python(DeltaTable API):

python
from delta.tables import *
deltaTablePeople = DeltaTable.forName(spark, "people10m")
dfUpdates = DeltaTable.forName(spark, "people10mupdates").toDF()

deltaTablePeople.alias('people') \
  .merge(dfUpdates.alias('updates'), 'people.id = updates.id') \
  .whenMatchedUpdate(set = {"id": "updates.id", "firstName": "updates.firstName" /* ... */}) \
  .whenNotMatchedInsert(values = {"id": "updates.id", "firstName": "updates.firstName" /* ... */}) \
  .execute()

条件付きアクション(AND 条件・複数 MATCHED 句)

sql
-- 一致のうち削除フラグ付きは削除、それ以外は更新(複数 MATCHED 句は順に評価、最後以外は条件必須)
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

-- 条件付き挿入(既定値も指定可)
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-only マージ・冪等)

sql
MERGE INTO logs
USING newDedupedLogs
ON logs.uniqueId = newDedupedLogs.uniqueId
WHEN NOT MATCHED
  THEN INSERT *

日付範囲で最適化(直近 7 日のみ探索):

sql
MERGE INTO logs
USING newDedupedLogs
ON logs.uniqueId = newDedupedLogs.uniqueId AND logs.date > current_date() - INTERVAL 7 DAYS
WHEN NOT MATCHED AND newDedupedLogs.date > current_date() - INTERVAL 7 DAYS
  THEN INSERT *

注意: 新データセット 内部の重複は防げない。マージ前に新データ側を重複排除しておくこと。

WHEN NOT MATCHED BY SOURCE(ソースにない行の更新/削除・段階同期)

sql
-- ソースにないターゲット行を削除(条件で範囲限定を推奨)
MERGE INTO target USING source
  ON target.key = source.key
  WHEN NOT MATCHED BY SOURCE THEN DELETE

-- 直近 5 日で増分同期: 更新・挿入・(範囲内の未一致を)削除
MERGE INTO target AS t
USING (SELECT * FROM source WHERE created_at >= (current_date() - INTERVAL '5' DAY)) AS s
ON t.key = s.key
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *
WHEN NOT MATCHED BY SOURCE AND created_at >= (current_date() - INTERVAL '5' DAY) THEN DELETE

WHEN NOT MATCHED BY SOURCE を条件なしで使うとターゲット全書き換えになり高コスト。必ず条件で絞る。

スキーマ進化つきマージ

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

SCD / CDC: Lakeflow パイプラインは SCD タイプ 1・2 の追跡/適用をネイティブ対応。CDC フィードは AUTO CDC ... INTO(順不同レコードを正しく処理)を使う。手書き MERGE で CDC を扱う場合は、キーごとに最新変更のみ残すよう ソースを前処理してから適用する(複数一致失敗の回避にもなる)。

4-6. ETL クイックスタート(Auto Loader → Delta、命令的変換)

Auto Loader で JSON を増分取り込み → Delta テーブルへ書き込み(bronze 相当)

python
from pyspark.sql.functions import col, current_timestamp

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"

spark.sql(f"DROP TABLE IF EXISTS {table_name}")
dbutils.fs.rm(checkpoint_path, True)

(spark.readStream
  .format("cloudFiles")                                   # Auto Loader
  .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))

ポイント:

  • format("cloudFiles") が Auto Loader。到着した新ファイルのみを増分検出・処理。
  • cloudFiles.schemaLocationcheckpointLocation により、スキーマ追跡と 正確な増分処理(再実行に強い=冪等的) を実現。
  • .select(...) が変換(メタデータ列 source_file・処理時刻 processing_time を付与)。
  • 既定のテーブル形式は Delta Lake(ACID トランザクション)。
  • ノートブックはジョブのタスクとしてスケジュールし、本番運用(オーケストレーション)へ載せる。

読み取り・確認:

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

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

  1. 3 アクションの正確な対応(最頻出)
    • warn=EXPECT / dp.expect(保持・件数記録、既定)、drop=ON VIOLATION DROP ROW / dp.expect_or_drop(削除・件数記録)、fail=ON VIOLATION FAIL UPDATE / dp.expect_or_fail(更新失敗・ロールバック・メトリックは記録されない)。
  2. fail 時のパイプライン挙動: トリガー型は当該フローのみロールバック(他フローは継続)、連続型はフロー+全依存フローが停止。
  3. 複数期待値のグルーピングは Python 限定expect_all / _or_drop / _or_fail、辞書引数)。SQL は複数 CONSTRAINT をカンマ区切りで並べるのみ。
  4. 制約に使えないもの: カスタム Python 関数・外部サービス呼び出し・他テーブル参照サブクエリ。期待値名はデータセット内で一意。
  5. 隔離(quarantine)パターン: データを捨てずに is_quarantined フラグ+一時テーブル/ビューで有効・無効を分離。品質メトリックも取れる。
  6. 検証テーブルはオーケストレーションではない: ダウンストリームは自動ゲートされない。止めたいなら別パイプライン+ジョブのタスク依存。
  7. MERGE の句ルール: 3 種類(MATCHED / NOT MATCHED [BY TARGET] / NOT MATCHED BY SOURCE)。同種の句が複数なら最後以外は条件必須。NOT MATCHED BY SOURCE はターゲット列のみ参照可・条件で範囲を絞らないと全書き換え。
  8. 複数一致は失敗: ソースの複数行がターゲット同一行に一致すると曖昧で失敗 → ソースを前処理して一意化(CDC は最新のみ残す)。
  9. 重複排除の冪等性: insert-only マージ(WHEN NOT MATCHED THEN INSERT *)で再挿入を防ぐ。ただし 新データ内部の重複は事前に排除。日付範囲条件で最適化。foreachBatch でストリーミング継続重複排除。
  10. スキーマ強制 vs 進化: Delta は既定で強制。MERGE WITH SCHEMA EVOLUTION(DBR 15.2+)で自動進化。UPDATE SET */INSERT * はソース列=ターゲット列前提。
  11. バッチ vs ストリーミング: バッチ=全再処理・単純・正確・低速、ストリーミング=新データのみ・高速・ステートフルで複雑。メダリオンは bronze=ストリーミング寄り、silver/gold=バッチ(具体化ビューの増分更新)寄り。
  12. データ品質の監視の 2 本柱: 異常検知(鮮度=コミット遅延、完全性=24h 行数)とデータプロファイル(分布・null 率・パーセンタイル・ドリフト、カスタムメトリック・時間粒度)。サーバーレスで実行しテーブルを変更しない。
  13. Expectations(インライン品質)と 監視(事後品質)の役割分担を説明できること。
  14. Auto Loader: format("cloudFiles")cloudFiles.formatcloudFiles.schemaLocationcheckpointLocationtrigger(availableNow=True) の役割。増分・スキーマ進化・冪等。

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

  • [ ] warn / drop / fail の SQL 構文(EXPECT / ON VIOLATION DROP ROW / ON VIOLATION FAIL UPDATE)と Python API(dp.expect / expect_or_drop / expect_or_fail)を対応付けて言える。
  • [ ] fail はメトリックが記録されず更新がアトミックにロールバックされること、トリガー型と連続型で停止範囲が違うことを説明できる。
  • [ ] 複数期待値の集合アクション(expect_all / _or_drop / _or_fail)が Python 限定で、辞書を引数に取ることを知っている。
  • [ ] 制約句に使えないもの(カスタム Python 関数・外部呼び出し・他テーブル参照サブクエリ)を挙げられる。
  • [ ] 無効レコードの隔離(quarantine)を is_quarantined フラグ+一時テーブル/ビューで実装できる。
  • [ ] 検証テーブルはオーケストレーションではなく、ダウンストリームを止めるには別パイプライン+ジョブ依存が必要だと理解している。
  • [ ] MERGE INTO の 3 種類の WHEN 句と、それぞれで使えるアクション(DELETE/UPDATE/INSERT)を区別できる。
  • [ ] 同種の WHEN 句が複数あるとき最後以外は条件必須、NOT MATCHED BY SOURCE はターゲット列のみ参照可・要範囲限定であることを言える。
  • [ ] 複数一致(multiple matches)が失敗する理由と、ソース前処理による回避策を説明できる。
  • [ ] insert-only マージによる冪等な重複排除ができ、「新データ内部の重複は事前排除が必要」を理解している。
  • [ ] MERGE WITH SCHEMA EVOLUTION とスキーマ強制の違い、UPDATE SET * / INSERT * の前提を説明できる。
  • [ ] 欠損値・型不一致・重複・スキーマ変化に対するクレンジング手段を具体例で挙げられる。
  • [ ] バッチとストリーミングのセマンティクス差(再処理・遅延データ・ステートフル)と、メダリオン各層の推奨を説明できる。
  • [ ] データ品質の監視の異常検知(鮮度・完全性)とデータプロファイル(分布・ドリフト・カスタムメトリック)の違いと、Expectations との役割分担を言える。
  • [ ] Auto Loader(cloudFiles)+ Delta + チェックポイントで増分・スキーマ進化・冪等な取り込みができる。