テーマ切替
Professional 学習教材 ⑨ データモデリング(配点 6%)
本教材は Databricks 認定データエンジニア Professional(上位資格) 試験の「データモデリング」ドメイン(配点 6%)向けの自習教材です。 SCD Type 1/2、AUTO CDC(旧 APPLY CHANGES)、AUTO CDC FROM SNAPSHOT、変更データフィード(CDF)、メダリオン/ディメンションモデルを、公式ドキュメント(日本語)を実読して深く扱います。この md 単体で学習が完結するように構成しています。
参照した公式ページ(すべて WebFetch で実読、2026-07 時点の日本語版):
- 変更データキャプチャとスナップショット(CDC/SCD 概念): https://learn.microsoft.com/ja-jp/azure/databricks/ldp/what-is-change-data-capture
- AUTO CDC API(旧 APPLY CHANGES): https://learn.microsoft.com/ja-jp/azure/databricks/ldp/cdc
- 変更データフィード(CDF)を使用する: https://learn.microsoft.com/ja-jp/azure/databricks/delta/delta-change-data-feed
- メダリオン レイクハウス アーキテクチャ: https://learn.microsoft.com/ja-jp/azure/databricks/lakehouse/medallion
- チュートリアル: CDC を使用して ETL パイプラインを構築する: https://learn.microsoft.com/ja-jp/azure/databricks/ldp/tutorial-pipelines
注: 一部ページはアクセス時に正規 URL(例:
data-engineering/what-is-cdc、tables/features/change-data-feed)へリダイレクトされますが、いずれも日本語版の同一内容を取得済みです。
1. このドメインの概要
「データモデリング」は、上流の運用システム(Oracle / PostgreSQL / SQL Server / MySQL など)の変更を、分析・レポート・機械学習のために Databricks 側のテーブルへ正確かつ最新の状態に同期し続けるための設計・実装を扱うドメインです。Professional 試験では、次の 4 つの柱を「実装レベルで」問われます。
- メダリオンアーキテクチャ(Bronze → Silver → Gold)とデータモデリングの位置づけ(正規化/非正規化、ディメンショナルモデル/スタースキーマ)。
- 変更データキャプチャ (CDC) の考え方と、ソース形式(変更フィード / スナップショット)に応じた処理方式の選択。
- AUTO CDC API(旧 APPLY CHANGES API) による SCD Type 1 / Type 2 の自動計算(
__START_AT/__END_AT、シーケンス列、順序外イベント処理、AUTO CDC FROM SNAPSHOT)。 - 変更データフィード (CDF) による行レベル変更の追跡と、下流への増分伝播(
readChangeFeed/table_changes())。
学習の核となる概念:
- 現在の状態だけ欲しいのか(SCD Type 1)、変更の全履歴が欲しいのか(SCD Type 2) をユースケースから選べること。
- ソースが変更フィードを出すか、スナップショットしか出さないか で
AUTO CDCとAUTO CDC FROM SNAPSHOTを使い分けられること。 - 順序外(out-of-order)イベントを、シーケンス列によって決定論的に正しく処理できる仕組みを説明できること。
- CDF(変更を書き込む/読み取る仕組み) と AUTO CDC(変更を適用して SCD を作る仕組み) の役割の違いを混同しないこと。
2. 重要用語集
| 用語(日本語) | English | 説明 |
|---|---|---|
| 変更データキャプチャ | Change Data Capture (CDC) | データベースを「静的な完全データ」ではなく「一連の変更(挿入/更新/削除)」として扱うデータ統合パターン。変更されたレコードのみを含むフィードを生成する。 |
| CDC フィード | CDC feed | 挿入・更新・削除のみを含む変更ストリーム。各行に操作種別(INSERT/UPDATE/DELETE)、データ値、決定論的順序付け用のシーケンス番号/タイムスタンプを持つ。 |
| スナップショット | Snapshot | 特定時点のテーブルの完全な状態(全行)。変更のみを含む CDC フィードと対比される。CDC が使えない場合の唯一の選択肢になりうる。 |
| 緩やかに変化するディメンション | Slowly Changing Dimension (SCD) | 上流の変更を分析テーブルに適用・モデル化する方式の総称。Type 1(現在の状態のみ)と Type 2(全履歴)が代表。 |
| SCD Type 1 | SCD Type 1 | 変更のたびに古い値を新しい値で上書きし、各レコードの最新バージョンのみ保持。履歴は残さない。 |
| SCD Type 2 | SCD Type 2 | 各バージョンを別行として保持し完全な履歴を残す。__START_AT/__END_AT で各バージョンの有効期間を表す。アクティブ行は __END_AT = NULL。 |
| AUTO CDC(旧 APPLY CHANGES) | AUTO CDC (formerly APPLY CHANGES) | Lakeflow パイプラインで CDC フィードから SCD Type 1/2 を自動計算する API。APPLY CHANGES を置き換える(構文は同一)。SQL の AUTO CDC ... INTO / Python の create_auto_cdc_flow()。 |
| AUTO CDC FROM SNAPSHOT | AUTO CDC FROM SNAPSHOT | 連続するスナップショットを比較して合成変更フィードを生成し、SCD Type 1/2 を適用する API。Python のみサポート。create_auto_cdc_from_snapshot_flow()。 |
| 変更データフィード | Change Data Feed (CDF) | Delta / Iceberg v3 テーブルのバージョン間の行レベル変更を追跡する仕組み。デルタ独自の CDC フィード。readChangeFeed / table_changes() で読む。 |
| シーケンス列 | Sequence column / SEQUENCE BY | イベントの正しい順序を決める列。単調増加で、各シーケンス値でキーごとに 1 更新。NULL 不可、並べ替え可能な型が必須。順序外イベントを正しく適用するために使う。 |
__START_AT / __END_AT | __START_AT / __END_AT | SCD Type 2 でレコードの各バージョンの有効期間を表す列。シーケンス値が伝播される。__END_AT = NULL は現在アクティブなレコード。 |
| ストリーミングテーブル | Streaming Table | 追加された新規データのみを増分処理するテーブル。AUTO CDC のターゲットに使う(CREATE STREAMING TABLE / create_streaming_table())。 |
| マテリアライズドビュー | Materialized View | クエリ結果を事前計算・保持するビュー。Gold 層の集計に使う(CREATE MATERIALIZED VIEW)。 |
| ファクト / ディメンション | Fact / Dimension | ファクト=測定値(売上など)を持つ中心テーブル。ディメンション=分析軸(顧客・製品・地域・時間など)。 |
| スタースキーマ | Star Schema | 中央のファクトテーブルを複数のディメンションテーブルが取り囲む次元モデル。Gold 層で用いられる。 |
| ディメンショナルモデル | Dimensional Model | ファクトとディメンションでビジネスをモデル化する手法。Gold 層でリレーションとメジャーを定義。 |
| メダリオンアーキテクチャ | Medallion Architecture | Bronze(生)→ Silver(検証済)→ Gold(エンリッチ済)の多層でデータ品質を段階的に高める設計パターン。別名マルチホップアーキテクチャ。 |
| 主キー / 外部キー | Primary Key / Foreign Key | 主キー=行を一意に識別する列(AUTO CDC の KEYS)。外部キー=別テーブルの主キーを参照する列。次元モデルでファクト⇔ディメンションを結ぶ。 |
| サロゲートキー | Surrogate Key | 業務的意味を持たない安定した代理キー。CDC により時間経過に対して安定して使える。 |
| 初期ハイドレーション | Initial hydration | 既存テーブル全体を最初に一度だけ読み込むこと。AUTO CDC の once(ワンスフロー)で一括ロード後、継続処理に移行。 |
| バイテンポラル追跡(ベータ) | Bitemporal tracking (Beta) | SCD Type 2 を拡張し「業務時間」と「システム時間」の 2 軸で履歴を追跡。__SYSTEM_START_AT/__SYSTEM_END_AT を追加。SYSTEM SEQUENCE BY。 |
_change_type | _change_type | CDF のメタデータ列。insert / update_preimage(更新前) / update_postimage(更新後) / delete。 |
_commit_version / _commit_timestamp | _commit_version / _commit_timestamp | CDF のメタデータ列。変更が含まれる Delta ログ/テーブルバージョン、コミットのタイムスタンプ。 |
| Lakeflow 宣言型パイプライン | Lakeflow Declarative Pipelines(旧 DLT / Delta Live Tables) | ストリーミングテーブル/マテリアライズドビューを宣言的に定義する ETL フレームワーク。AUTO CDC はこの上で動く。 |
| 期待値(データ品質制約) | Expectations | パイプラインでデータ品質を検証する制約。EXPECT ... ON VIOLATION DROP ROW など。 |
| 自動ローダー | Auto Loader | クラウドストレージから増分取り込み。スキーマ推論・進化を自動処理(cloudFiles / read_files)。 |
3. 詳細解説
3-1. メダリオンとデータモデリングの位置づけ(ディメンショナルモデル、正規化/非正規化)
メダリオンアーキテクチャは、レイクハウス内でデータ品質を層ごとに段階的に高めるデータ設計パターンです(別名「マルチホップアーキテクチャ」)。Bronze ⇒ Silver ⇒ Gold と流れるにつれ、データの構造と品質が向上します。推奨されるが必須ではない点に注意。
| 観点 | Bronze(銅・生) | Silver(銀・検証済) | Gold(金・エンリッチ済) |
|---|---|---|---|
| この層で起こること | 生データの取り込み | データのクレンジングと検証 | ディメンショナルモデリングと集計 |
| 主な利用者 | データエンジニア、データ運用、コンプライアンス/監査 | データエンジニア、データアナリスト、データサイエンティスト | BI 開発者/ビジネスアナリスト、DS/ML エンジニア、経営層、運用チーム |
Bronze 層(生データ取り込み)
- ソースの生の状態を元の形式のまま保存・維持。増分追加され時間とともに増大。
- 直接アナリストが使うのではなく、Silver を作るワークロードが使う。単一の信頼できるソースとして機能し、全履歴を保持して再処理・監査を可能にする。
- クラウドオブジェクトストレージ(S3/GCS/ADLS)、メッセージバス(Kafka/Kinesis)、Lakehouse フェデレーションなどストリーミング/バッチを任意に組み合わせ。
- データ検証は最小限。削除データを確実に防ぎスキーマ変更から守るため、多くのフィールドを string / VARIANT / binary で保存することが推奨。
_metadata.file_nameなどメタデータ列を付与することも。
Silver 層(検証・重複除去)
- 1 つ以上の Bronze/Silver テーブルから読み、Silver へ書く。取り込みから直接 Silver に書くのは非推奨(スキーマ変更や破損レコードでエラーになるため)。追加専用前提で Bronze からの読み取りはストリーミング読み取りを基本にし、バッチ読み取りは小さなディメンションテーブルに限定。
- 常に「各レコードの検証済み・非集計表現」を少なくとも 1 つ含む。
- 実施内容: スキーマ適用、null/欠損値の処理、重複除去、順序外・遅延到着データの解決、データ品質チェック/適用、スキーマ展開、型キャスト、結合。
- ここでデータモデリングを開始(VARIANT、JSON 文字列、struct/map/array、フラット化 or 複数テーブルへの正規化)。
Gold 層(分析エンリッチ)
- 集計・フィルタ済みの高度に洗練されたビュー。ビジネス機能に対応した意味的に意味のあるデータセット。
- ディメンションモデルでレポート/分析用にモデル化し、リレーションを確立しメジャーを定義。人事・財務・IT など目的別に複数の Gold 層を作る顧客もいる。
- 集計(平均/カウント/最大/最小など)をマテリアライズドビューで事前計算(例:
weekly_sales)。頻繁に照会されるためパフォーマンス最適化を推奨。大量の履歴データは通常 Silver でアクセスし Gold には具体化しない。
正規化 vs 非正規化 / 次元モデル
- 正規化(重複排除・複数テーブル分割)は主に Silver で行い、データ品質と整合性を担保。
- 非正規化・ディメンショナルモデル(スタースキーマ) は Gold で行い、クエリ/ダッシュボードのパフォーマンスを優先。ファクト(測定値)とディメンション(分析軸)で構成する。
- SCD はこの次元モデルの実装技法。ディメンションテーブルが時間とともに変化する際、Type 1(上書き)か Type 2(履歴保持)を選ぶ。CDC により時間経過に対して安定したサロゲートキーが使える。
取り込み頻度とコスト(Gold のコスト制御): 継続的増分(高コスト/低レイテンシ)<トリガー増分<バッチ(低コスト/高レイテンシ)。宣言型では spark.readStream によるストリーミングテーブル、Trigger.AvailableNow などで制御。
3-2. 変更データキャプチャ(CDC)と AUTO CDC / AUTO CDC FROM SNAPSHOT
CDC の考え方: データベースを完全な静的データではなく「一連の変更」として扱う。従業員テーブルの 1 行が更新されると、CDC フィードにはその 1 行の UPDATE レコードだけが生成される。各行には操作種別と、順序を決定論的に並べ替える列(例: sequenceNum)が付く。
CDC の利点:
- 変更データは完全データより小さく、下流クエリで増分更新として処理できる(マテリアライズドビューを増分リフレッシュ可能)。
- 特定時点のレコードを再構築できる形で保存でき、監査・ポイントインタイムレポート・傾向分析の完全履歴を提供。
- 時間経過に対して安定したサロゲートキーが使える。
ソース形式の課題: ソースは形式がまちまち。個々の変更(挿入/更新/削除)を出すフィードもあれば、テーブル全体の定期スナップショットしか出さないものもある。形式ごとに処理方式が異なる。従来はカスタム MERGE INTO(ステージングテーブル、ウィンドウ関数、シーケンス前提が必要で複雑・エラーが起きやすい)に依存していた。
CDC フィードとは: ソースの変更のみを含むフィード。各 CDC レコードは、①操作種別(INSERT/UPDATE/DELETE)、②レコードのデータ値、③決定論的順序付け用のシーケンス番号/タイムスタンプ、を持つ。シーケンス番号により遅延到着・順序外到着を正しく適用できる。SQL Server/MySQL/Oracle はネイティブに CDC フィードを生成。Delta テーブルは独自の CDC フィード=CDF を生成する。
スナップショットとは: 特定時点のテーブルの完全な状態(全行)。CDC フィードが常に有効とは限らない理由: コスト(運用 DB の負荷増)、ソース DB のパフォーマンス、CDC 非対応のレガシー、組織的制約(上流 DB を所有していない)。スナップショットの取得元: RDB の定期エクスポート、クラウドストレージのファイルダンプ、差分テーブル(各バージョンが実質スナップショット)、Delta Sharing。スナップショットはレコードレベルの変更をキャプチャしないため、スナップショット間を比較して挿入/更新/削除を推測する必要がある。
API の使い分け(重要):
| 条件 | 使う API | インターフェイス |
|---|---|---|
| ソースが変更フィードを出す(CDF 有効な Delta、Debezium/Oracle GoldenGate 等の RDB CDC、シーケンス列付きの挿入/更新/削除ストリーム) | AUTO CDC | SQL + Python |
| ソースが CDC 非対応で定期スナップショット(全テーブルダンプ)のみ | AUTO CDC FROM SNAPSHOT | Python のみ |
どちらでも、現在の状態のみ必要なら SCD Type 1、監査/ポイントインタイム/傾向分析で全履歴が必要なら SCD Type 2 を選ぶ。
AUTO CDC のしくみ:
- ストリーミングテーブルを作り、
AUTO CDC ... INTO(SQL)/create_auto_cdc_flow()(Python)で、変更フィードのソース・キー・シーケンスを指定。 - シーケンス列で定義した順序でイベントを処理し、順序外レコードを自動処理する。シーケンス列は「各シーケンス値でキーごとに 1 更新」となる単調増加の表現である必要があり、
NULLは不可。SCD Type 2 ではシーケンス値がターゲットの__START_AT/__END_ATに伝播される。 - 初期ハイドレーション: 既存 DB テーブル全体を先に読み込む必要がある場合、ワンスフロー(
once) モードで利用可能な全データを一度だけ処理して停止。その後トリガー/連続モードで継続 CDC。一括・増分で一貫したロジックになる。
AUTO CDC FROM SNAPSHOT のしくみ:
- 連続するスナップショットを比較して変更を判定。差分テーブル/クラウドストレージファイル/JDBC から直接スナップショットを読める。Python のみ。
- 処理内容: ①連続スナップショットを比較し挿入/更新/削除を識別 → ②差分から合成変更フィードを生成 → ③
AUTO CDCと同じ SCD Type 1/2 ロジックを適用。 - 初期読み込みだけではない。スナップショットが唯一の形式である場合の継続処理用。新しいスナップショット到着のたびに前回と比較して変更を導出。
- 重要な制限: あるスナップショットから次への変更のみ認識し、中間変更は取得しない(例: 日次スナップショットで 1 日に住所を A→B→C と 2 回変えると、変更フィードは A→C に直接飛ぶ)。
- スナップショットはバージョン昇順で処理する必要があり、順序外のスナップショットは無視される。
AUTO CDC FROM SNAPSHOT の 2 つの処理パターン:
- パイプライン取り込み時間を使う: 実行時にスナップショットを読み、取り込み時間をバージョンとする。更新ごとに新スナップショット取り込み。連続モードではトリガー間隔で複数取り込み。→ スナップショットが定期的・順番に到着する場合に使用。
- バージョン関数を使う: 処理すべきスナップショットを返す関数を指定。関数は
(DataFrame, version_number)のタプルを返し、API はバージョン番号の昇順で処理。順序外は無視、新規なしならNone。→ 複数同時到着・順序外到着・順序を明示制御したい場合に使用。
その他の AUTO CDC 機能:
- ターゲットへの DML: 標準ストリーミングテーブルと異なり、AUTO CDC ターゲット(Unity Catalog テーブル)はパイプライン実行中でも
INSERT/UPDATE/DELETE/MERGEをサポート。 - ターゲットからの CDF 読み取り: AUTO CDC ターゲットは独自の CDF を出力でき、下流パイプラインが AUTO CDC 出力の変更を消費できる。
- メトリクス: 実行ごとに
num_upserted_rows/num_deleted_rowsを自動キャプチャ。 - 要件: CDC API は Pro または Advanced エディション(あるいはサーバーレス)のパイプラインが必要。
AUTO CDCは Apache Spark 宣言パイプラインでは非サポート。
バイテンポラル追跡(ベータ、AUTO CDC のみ): SCD Type 2 を拡張し「業務時間(イベントが実際に起きた時刻)」と「システム時間(システムが記録/取り込んだ時刻)」の 2 軸を追跡。STORED AS BITEMPORAL(SQL)/ stored_as_scd_type="bitemporal"(Python)、業務時間に SEQUENCE BY、システム時間に SYSTEM SEQUENCE BY。ターゲットに __SYSTEM_START_AT/__SYSTEM_END_AT が追加される。順序外イベントが来ると末尾追加ではなく影響を受ける履歴を修正する。
3-3. SCD Type 1 と Type 2(違い・計算・__START_AT/__END_AT・順序付け)
Type 1(現在の状態のみ): 変更のたびに古い値を新しい値で上書き、最新バージョンのみ保持。履歴なし。使いどころ = 現在の状態のみ必要 / 下流マテビューを増分リフレッシュ / 結合に安定サロゲートキーが必要。「最終テーブルのみを保存」する簡潔な方式。Owner → Manager に変われば Manager のみ残る。
Type 2(履歴追跡): メタデータでタイムスタンプ付き複数バージョンを保持し、完全な履歴を残す。__START_AT/__END_AT が各バージョンの有効期間を定義、アクティブ行は __END_AT = NULL。使いどころ = 監査/規制で履歴が必要 / エンティティの進化を理解 / ポイントインタイムレポート / 傾向分析。
順序付け(順序外イベント処理)の要点:
- シーケンス列は単調増加、キーごとに各シーケンス値で 1 更新、
NULL不可、並べ替え可能な型。 - 遅延/順序外イベントもシーケンス値で正しい最終状態に収束する。
- 複数列での順序付け: タイムスタンプ + ID でタイブレークするなど。
STRUCTで結合し、最初のフィールド優先、同値なら次のフィールドで判定。- SQL:
SEQUENCE BY STRUCT(timestamp_col, id_col) - Python:
sequence_by = struct("timestamp_col", "id_col")
- SQL:
具体例(公式データセット) — 入力 CDC(sequenceNum が順序):
| userId | name | city | operation | sequenceNum |
|---|---|---|---|---|
| 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 |
Type 1 の結果(最新のみ、削除は消える、seq=5 の Mercedes は seq=6 により破棄):
| userId | name | city |
|---|---|---|
| 124 | Raul | Oaxaca |
| 125 | Mercedes | Guadalajara |
| 126 | Lily | Cancun |
(TRUNCATE(sequenceNum=3)を有効にすると、seq=3 でテーブルがクリアされ、それ以前に取り込まれた 124/126 が消え、最終的に 125 Guadalajara のみ残る。TRUNCATE は sequenceNum 順で「その時点までをクリア」と考える。)
Type 2 の結果(全履歴、__END_AT = null が現在アクティブ):
| userId | name | city | __START_AT | __END_AT |
|---|---|---|---|---|
| 123 | Isabel | Monterrey | 1 | 5 |
| 123 | Isabel | Chihuahua | 5 | 6 |
| 124 | Raul | Oaxaca | 1 | null |
| 125 | Mercedes | Tijuana | 2 | 5 |
| 125 | Mercedes | Mexicali | 5 | 6 |
| 125 | Mercedes | Guadalajara | 6 | null |
| 126 | Lily | Cancun | 2 | null |
ユーザー 123 は削除で seq=6 にて終了(2 バージョン)。125 は都市変更で 3 バージョン。シーケンス値が __START_AT/__END_AT にそのまま入っている点に注目。
Type 2 で列サブセットのみ追跡: 既定ではどの列が変わっても新バージョンを作るが、TRACK HISTORY ON(SQL)/ track_history_except_column_list(Python)で追跡する列を限定できる。追跡対象外の列の変更は新履歴を作らず現在バージョンを上書き。重要属性の履歴を残しつつストレージ/クエリ複雑度を削減。上例で city を除外すると 123 は 1 行(Chihuahua, 1→6)、125 は 1 行(Guadalajara, 2→null)に集約される。
削除・切り捨ての扱い(オプション、どちらも省略可):
APPLY AS DELETE WHEN/apply_as_deletes: 条件に合う行を削除として扱う。Type 2 では該当バージョンを__END_ATで終了。APPLY AS TRUNCATE WHEN/apply_as_truncates: そのシーケンス時点でテーブルをクリア(Type 1 系の全消去)。
3-4. 変更データフィード(CDF)による下流への増分伝播
CDF とは: Delta Lake / Apache Iceberg v3 テーブルのバージョン間の行レベル変更を追跡する仕組み。デルタ独自の CDC フィード。用途 = 増分 ETL(前回実行以降に変わった行だけ処理)、監査証跡、データレプリケーション(下流テーブル/キャッシュ/外部システムへの同期)。
Databricks は 2 方式をサポート:
- 自動変更データフィード(パブリックプレビュー): 行系列メタデータを使い、書き込み時ではなくクエリ(読み取り)時に変更を計算。個別のテーブル構成不要。Delta Lake と Iceberg v3 で動作。
MERGE/UPDATE書き込みごとに変更を計算しないため、レガシーより書き込み性能が高くストレージコストが低い。要件: Databricks Runtime 18 以降、Unity Catalog 登録、行追跡有効。 - レガシー変更データフィード: テーブルへの書き込み時に変更を記録。Delta Lake のみ。個別テーブルで明示的に有効化が必要。
両方式は同じ読み取り API(readChangeFeed / table_changes())を使い、バッチ/構造化ストリーミング/Delta Sharing で動作。レガシーと自動は併用不可。
メタデータ列(CDF 読み取り時に付与):
| 列名 | 型 | 値 |
|---|---|---|
_change_type | String | insert / update_preimage(更新前値) / update_postimage(更新後値) / delete |
_commit_version | Long | その変更を含む Delta ログ/テーブルバージョン |
_commit_timestamp | Timestamp | コミット作成時のタイムスタンプ |
スキーマにこれらと同名の列があると CDF は使えない(有効化前に列名変更で解消)。
レガシー CDF の有効化(テーブルプロパティ delta.enableChangeDataFeed):
- 新規:
CREATE TABLE ... TBLPROPERTIES (delta.enableChangeDataFeed = true) - 既存:
ALTER TABLE ... SET TBLPROPERTIES (delta.enableChangeDataFeed = true) - 一時的に無効→再有効化すると、その間隔はクエリ不可(間隔中の変更を照会するには自動 CDF を使う)。
読み取り(バッチ):
- SQL:
SELECT * FROM table_changes('tableName', 開始バージョン[, 終了バージョン])(開始/終了は含む、開始のみ指定で最新まで)。タイムスタンプはyyyy-MM-dd[ HH:mm:ss[.SSS]]文字列。 - Python/Scala:
.option("readChangeFeed","true").option("startingVersion", 0).option("endingVersion", 10)またはstartingTimestamp/endingTimestamp。 - バッチ読み取りには開始バージョンが必須。CDF 有効化前のバージョンを指定するとエラー。
読み取り(ストリーミング/増分処理・推奨):
- Databricks は構造化ストリーミング + CDF による増分処理を推奨(バージョンを自動追跡)。
- ストリーム初回起動時、CDF はテーブルの最新スナップショットを
INSERTレコードとして返し、以降を変更データとして返す。 - ターゲットに既に特定時点までの変更が入っている場合は、
startingVersionを指定して既存状態をINSERTとして再処理しないようにする(チェックポイント破損からの復旧など)。 - レート制限:
maxFilesPerTrigger/maxBytesPerTrigger/excludeRegexをサポート。開始スナップショット以外ではコミット全体にアトミックに適用。
保持・アーカイブ・制限:
- CDF は恒久的な変更記録ではない。CDF 有効化後に発生した変更のみを記録。レコードは一時的で、指定保持期間内のみアクセス可。バージョン削除後はその CDF を読めない(マネージドテーブルは履歴を自動クリーンアップするので、指定した開始バージョンは最終的に削除される)。
- 恒久履歴が必要なら、
trigger.AvailableNowなどで CDF から新テーブルへ増分書き込みしてアーカイブ。 - 範囲外バージョン: 既定では最後のコミットを超える指定は
timestampGreaterThanLatestCommitエラー。SET spark.databricks.delta.changeDataFeed.timestampOutOfRange.enabled = trueで、開始超過→空結果 / 終了超過→最後まで返す。 - 列マッピング/非加法スキーマ変更の制限: 列の名称変更・削除・型変更・NULL 許容変更(非加法変更)を跨ぐ範囲は読めない。範囲の終了バージョンのスキーマを使用。跨ぐ場合は範囲を分割する。
- 自動 CDF の追加制限: 複数ステートメントトランザクション中にソースが変わると非サポート、行フィルター/列マスクのテーブルは非サポート、外部 Iceberg リーダーからはクエリ不可(Databricks リーダーのみ)。
CDF と AUTO CDC の役割の違い(頻出):
- CDF = テーブルの行レベル変更を記録・読み取る低レベル機能(変更を出す/取り込む)。
- AUTO CDC = 変更を受けて SCD Type 1/2 を適用しテーブルを再具体化する高レベル API(変更を適用して整合したディメンションを作る)。
- チュートリアルの記述: 「Delta は CDF をサポートし
table_changesで変更を照会できるが、CDF の主用途はパイプライン内の変更キャプチャであり、最初からテーブル変更の完全ビューを作ることではない。順序外イベントがあると実装が複雑。Lakeflow パイプライン(AUTO CDC)はこの複雑さを解消し、_sequence_byに基づき順序外レコードを自動処理する」。
4. 構文・コード例
4-1. AUTO CDC INTO — SCD Type 1(SQL / Python)
SQL:
sql
CREATE OR REFRESH STREAMING TABLE users_current;
CREATE FLOW apply_cdc AS AUTO CDC INTO
users_current
FROM
stream(main.cdc_tutorial.users_cdf)
KEYS
(userId)
APPLY AS DELETE WHEN
operation = "DELETE"
APPLY AS TRUNCATE WHEN
operation = "TRUNCATE"
SEQUENCE BY
sequenceNum
COLUMNS * EXCEPT
(operation, sequenceNum)
STORED AS
SCD TYPE 1;Python:
python
from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr
@dp.view
def users():
return spark.readStream.table("main.cdc_tutorial.users_cdf")
dp.create_streaming_table("users_current")
dp.create_auto_cdc_flow(
target = "users_current",
source = "users",
keys = ["userId"],
sequence_by = col("sequenceNum"),
apply_as_deletes = expr("operation = 'DELETE'"),
apply_as_truncates = expr("operation = 'TRUNCATE'"),
except_column_list = ["operation", "sequenceNum"],
stored_as_scd_type = 1
)構文要素の対応表:
| SQL 句 | Python 引数 | 意味 |
|---|---|---|
AUTO CDC INTO <target> | target= | 適用先ストリーミングテーブル |
FROM stream(<source>) | source= | 変更フィードのソース |
KEYS (...) | keys=[...] | 一致に使う主キー |
SEQUENCE BY <col> | sequence_by= | 順序付け列(NULL 不可、並べ替え可能型) |
APPLY AS DELETE WHEN <cond> | apply_as_deletes= | 削除として扱う条件(任意) |
APPLY AS TRUNCATE WHEN <cond> | apply_as_truncates= | 切り捨てとして扱う条件(任意) |
COLUMNS * EXCEPT (...) / COLUMNS (...) | except_column_list= / column_list= | ターゲットに含める/除外する列 |
| `STORED AS SCD TYPE 1 | 2` | `stored_as_scd_type=1 |
| `TRACK HISTORY ON {cols | * EXCEPT(cols)}` | track_history_column_list= / track_history_except_column_list= |
| (なし) | ignore_null_updates= | null の更新を無視するか |
4-2. AUTO CDC INTO — SCD Type 2(SQL / Python)
SQL:
sql
CREATE OR REFRESH STREAMING TABLE users_history;
CREATE FLOW apply_cdc AS AUTO CDC INTO
users_history
FROM
stream(main.cdc_tutorial.users_cdf)
KEYS
(userId)
APPLY AS DELETE WHEN
operation = "DELETE"
SEQUENCE BY
sequenceNum
COLUMNS * EXCEPT
(operation, sequenceNum)
STORED AS
SCD TYPE 2;Python:
python
dp.create_streaming_table("users_history")
dp.create_auto_cdc_flow(
target = "users_history",
source = "users",
keys = ["userId"],
sequence_by = col("sequenceNum"),
apply_as_deletes = expr("operation = 'DELETE'"),
except_column_list = ["operation", "sequenceNum"],
stored_as_scd_type = "2"
)4-3. SCD Type 2 で列サブセットのみ追跡
SQL: 上記 Type 2 に TRACK HISTORY ON * EXCEPT (city) を追加。
sql
...
STORED AS
SCD TYPE 2
TRACK HISTORY ON * EXCEPT
(city)Python: create_auto_cdc_flow(..., stored_as_scd_type="2", track_history_except_column_list=["city"])
4-4. AUTO CDC FROM SNAPSHOT(Python のみ)
パターン A: 取り込み時間ベース(差分テーブル/クラウドストレージ/JDBC からソースを読む):
python
from pyspark import pipelines as dp
@dp.view(name="source")
def source():
return spark.read.table("main.cdc_tutorial.snapshot") # または csv / jdbc
dp.create_streaming_table("target")
dp.create_auto_cdc_from_snapshot_flow(
target = "target",
source = "source",
keys = ["userId"],
stored_as_scd_type = 2 # 1 にすると __START_AT/__END_AT 列は付かない
)初回実行後は全行がアクティブ(__START_AT=0, __END_AT=null)。新スナップショット到着後に再実行すると、削除→__END_AT 設定、更新→旧行終了+新行追加、挿入→新行、が自動生成される。
パターン B: バージョン関数ベース(順序を明示制御):
python
from typing import Optional, Tuple
from pyspark.sql import DataFrame
def next_snapshot_and_version(latest_snapshot_version: Optional[int]) -> Optional[Tuple[DataFrame, int]]:
# 次に処理すべきスナップショットを (DataFrame, version) で返す。
# 新規がなければ None を返す。順序外バージョンはスキップされる。
...
return (df, next_version)
dp.create_streaming_table("main.cdc_tutorial.target_versioned")
dp.create_auto_cdc_from_snapshot_flow(
target = "main.cdc_tutorial.target_versioned",
source = next_snapshot_and_version, # 関数を渡す
keys = ["userId"],
stored_as_scd_type = 2
)API はバージョン番号の昇順で処理。snapshot_3 の後に snapshot_4 を処理済みなら後から来た snapshot_3 はスキップ。
4-5. バイテンポラル AUTO CDC(ベータ)
SQL:
sql
CREATE OR REFRESH STREAMING TABLE target_bitemporal_sql;
CREATE FLOW target_bitemporal_sql AS AUTO CDC INTO
target_bitemporal_sql
FROM stream(cdc_source_sql)
KEYS (id)
SEQUENCE BY bt -- 業務時間
SYSTEM SEQUENCE BY st -- システム時間
STORED AS BITEMPORAL;Python: create_auto_cdc_flow(..., sequence_by="bt", system_sequence_by="st", stored_as_scd_type="bitemporal")。ターゲットに __START_AT/__END_AT(業務時間)と __SYSTEM_START_AT/__SYSTEM_END_AT(システム時間)が生成される。
4-6. CDF の有効化と読み取り
レガシー CDF 有効化:
sql
-- 新規
CREATE TABLE student (id INT, name STRING, age INT)
TBLPROPERTIES (delta.enableChangeDataFeed = true);
-- 既存
ALTER TABLE myDeltaTable SET TBLPROPERTIES (delta.enableChangeDataFeed = true);
-- レガシー→自動 CDF へ移行(レガシーを解除)
ALTER TABLE <table_name> UNSET TBLPROPERTIES ('delta.enableChangeDataFeed');バッチ読み取り(SQL):
sql
SELECT * FROM table_changes('tableName', 0, 10); -- v0〜v10
SELECT * FROM table_changes('tableName', '2021-04-21 05:45:46', '2021-05-21 12:00:00'); -- タイムスタンプ範囲
SELECT * FROM table_changes('tableName', 0); -- v0〜最新バッチ読み取り(Python):
python
spark.read \
.option("readChangeFeed", "true") \
.option("startingVersion", 0) \
.option("endingVersion", 10) \
.table("myDeltaTable")ストリーミング読み取り(増分処理・推奨):
python
(spark.readStream
.option("readChangeFeed", "true")
.table("myTable"))チェックポイント破損からの復旧(開始バージョン指定 + 新チェックポイント):
python
(spark.readStream
.option("readChangeFeed", "true")
.option("startingVersion", 76) # 既存ターゲットが v75 まで処理済み
.table("source_table")
.writeStream
.option("checkpointLocation", "<new-checkpoint-path>")
.toTable("target_table"))恒久履歴のアーカイブ(AvailableNow でバッチ化):
python
(spark.readStream
.option("readChangeFeed", "true")
.table("source_table")
.writeStream
.option("checkpointLocation", "<checkpoint-path>")
.trigger(availableNow=True)
.toTable("target_table"))4-7. メダリオン ETL パイプライン(チュートリアル要点)
Bronze(自動ローダー取り込み)→ Silver(期待値で品質検証)→ Gold(AUTO CDC で SCD1 具体化 / SCD2 履歴 + マテビュー集計)。
Bronze: 自動ローダー取り込み(SQL):
sql
CREATE OR REFRESH STREAMING TABLE customers_cdc_bronze
COMMENT "New customer data incrementally ingested from cloud object storage landing zone";
CREATE FLOW customers_bronze_ingest_flow AS
INSERT INTO customers_cdc_bronze BY NAME
SELECT * FROM STREAM read_files(
"/Volumes/<catalog>/<schema>/raw_data/customers",
format => "json",
inferColumnTypes => "true"
);Silver: 期待値(データ品質制約)(SQL):
sql
CREATE OR REFRESH STREAMING TABLE customers_cdc_clean (
CONSTRAINT no_rescued_data EXPECT (_rescued_data IS NULL) ON VIOLATION DROP ROW,
CONSTRAINT valid_id EXPECT (id IS NOT NULL) ON VIOLATION DROP ROW,
CONSTRAINT valid_operation EXPECT (operation IN ('APPEND','DELETE','UPDATE')) ON VIOLATION DROP ROW
);
CREATE FLOW customers_cdc_clean_flow AS
INSERT INTO customers_cdc_clean BY NAME
SELECT * FROM STREAM customers_cdc_bronze;(Python では expect_all_or_drop = {...} を create_streaming_table に指定。id は upsert に使うため null 不可。)
Gold: AUTO CDC で customers を具体化(SCD Type 1, SQL):
sql
CREATE OR REFRESH STREAMING TABLE customers;
CREATE FLOW customers_cdc_flow
AS AUTO CDC INTO customers
FROM stream(customers_cdc_clean)
KEYS (id)
APPLY AS DELETE WHEN operation = "DELETE"
SEQUENCE BY operation_date
COLUMNS * EXCEPT (operation, operation_date, _rescued_data)
STORED AS SCD TYPE 1;Gold: 履歴を SCD Type 2 で保持(SQL): 上記の customers を customers_history に、STORED AS SCD TYPE 2 に変えるだけ。
Gold: マテリアライズドビューで集計(SQL):
sql
CREATE OR REPLACE MATERIALIZED VIEW customers_history_agg AS
SELECT id,
count(distinct address) AS address_count,
count(distinct email) AS email_count,
count(distinct firstname) AS firstname_count,
count(distinct lastname) AS lastname_count
FROM customers_history
GROUP BY id;(SCD2 で重複除去済みのため、id ごとに直接カウントできる。)
参考: 週次売上のマテビュー(メダリオン Gold の例):
sql
CREATE OR REPLACE MATERIALIZED VIEW weekly_sales AS
SELECT week, prod_id, region,
SUM(units) AS total_units,
SUM(units * rate) AS total_sales
FROM orders
GROUP BY week, prod_id, region;5. 試験で問われるポイント
- SCD Type 1 と Type 2 の違いを即答できる: Type 1 = 上書き・最新のみ・履歴なし。Type 2 = 全履歴・
__START_AT/__END_AT・アクティブ行は__END_AT = NULL。 - どちらの API を使うかをソース形式から判断: 変更フィードあり →
AUTO CDC。スナップショットのみ →AUTO CDC FROM SNAPSHOT(Python のみ、SQL 非対応)。 - AUTO CDC は APPLY CHANGES の後継で構文は同一。旧称
APPLY CHANGES/create_apply_changes_flowも動くが AUTO CDC 推奨。 - シーケンス列の要件: 単調増加・並べ替え可能な型・
NULL不可。順序外イベントを正しく最終状態に収束させる。複数列はSTRUCT(SQL)/struct()(Python)で。 __START_AT/__END_ATの値: SCD2 でシーケンス値が伝播される。ターゲットが SCD Type 1 のときはこれらの列は付かない。- 順序外の Type 1/Type 2 の結果を読める(公式例:
seq=5の更新がseq=6により Type 1 では破棄、Type 2 では 5→6 のバージョンとして残る)。 TRACK HISTORY ON ... EXCEPT (...)で Type 2 の追跡列を限定 = 対象外列の変更は新履歴を作らず上書き。- AUTO CDC FROM SNAPSHOT の制限: 中間変更を取得しない(スナップショット間の差分のみ)。バージョン昇順で処理、順序外スナップショットは無視。2 パターン(取り込み時間ベース / バージョン関数ベース)の使い分け。
- CDF と AUTO CDC の役割の違い: CDF = 変更の記録/読み取り。AUTO CDC = 変更適用による SCD 具体化。混同しない。
- CDF の有効化: レガシーは
delta.enableChangeDataFeed = true(テーブルプロパティ)。自動 CDF は構成不要(RT18+・行追跡)。両者は併用不可。 - CDF のメタデータ列:
_change_type(insert/update_preimage/update_postimage/delete)、_commit_version、_commit_timestamp。同名列があると使えない。 - CDF 読み取り API:
table_changes('t', start[, end])(SQL、開始必須・両端含む)、readChangeFeed=true+startingVersion/endingVersionorstartingTimestamp/endingTimestamp。ストリーミング初回は最新スナップショットをINSERTとして返す。 - CDF は恒久記録ではない: 有効化後の変更のみ、保持期間内のみ、バージョン削除後は読めない。恒久履歴は別テーブルへアーカイブ。
- CDF の制限: 非加法スキーマ変更(列名変更/削除/型変更/NULL 許容変更)を跨ぐ範囲は読めない → 範囲分割。範囲外バージョンは既定エラー(
timestampOutOfRange.enabledで緩和)。 - メダリオンの各層の役割: Bronze=生(検証最小、string/VARIANT/binary 推奨、全履歴保持)、Silver=クレンジング/重複除去/正規化/検証(取り込みから直接書かない、ストリーミング読み取り基本)、Gold=ディメンショナルモデル/集計/マテビュー(BI 最適化)。
- ディメンショナルモデリングは Gold 層で行う。SCD はディメンションの実装技法。CDC は安定サロゲートキーを可能にする。
- CDC API の要件: Pro / Advanced エディション(またはサーバーレス)のパイプライン。
AUTO CDCは Apache Spark 宣言パイプラインでは非サポート。 - ストリーミングテーブル vs マテリアライズドビュー: AUTO CDC のターゲットはストリーミングテーブル。Gold の集計はマテビュー。
- 初期ハイドレーションは
once(ワンスフロー)で一括ロード後、継続 CDC に移行。
6. 理解度チェックリスト
- [ ] SCD Type 1 と Type 2 の違いを、上書き/履歴・
__END_AT = NULLの意味を含めて説明できる。 - [ ] ソースが変更フィードか、スナップショットのみかで
AUTO CDCとAUTO CDC FROM SNAPSHOTを選べる。 - [ ]
AUTO CDC FROM SNAPSHOTが Python のみ・中間変更を取得しない・バージョン昇順処理である点を説明できる。 - [ ]
AUTO CDCがAPPLY CHANGESの後継で構文が同一であることを知っている。 - [ ] シーケンス列の要件(単調増加/並べ替え可能型/NULL 不可)と、順序外イベントがどう処理されるかを説明できる。
- [ ] 公式の順序外サンプルで Type 1 / Type 2 の最終テーブルを手で導出できる(
seq=5がseq=6に負ける/両方が履歴に残る)。 - [ ]
AUTO CDC INTOの SQL(KEYS / SEQUENCE BY / APPLY AS DELETE WHEN / COLUMNS * EXCEPT / STORED AS SCD TYPE)を書ける。 - [ ]
create_auto_cdc_flow()の Python 引数(keys / sequence_by / apply_as_deletes / except_column_list / stored_as_scd_type)を書ける。 - [ ]
TRACK HISTORY ON ... EXCEPT (...)で Type 2 の追跡列を限定した場合の挙動を説明できる。 - [ ] CDF(変更の記録/読み取り)と AUTO CDC(SCD 具体化)の役割の違いを混同せず説明できる。
- [ ] レガシー CDF を
delta.enableChangeDataFeed = trueで有効化し、自動 CDF との違い(構成不要/読み取り時計算/併用不可)を説明できる。 - [ ] CDF のメタデータ列
_change_type(4 値)・_commit_version・_commit_timestampを挙げられる。 - [ ]
table_changes()とreadChangeFeedによるバッチ/ストリーミング読み取りを書け、開始/終了バージョンの意味を説明できる。 - [ ] CDF が恒久記録でないこと、保持期間・バージョン削除の影響、恒久履歴アーカイブ方法を説明できる。
- [ ] CDF の非加法スキーマ変更の制限(範囲分割が必要)を説明できる。
- [ ] メダリオン Bronze/Silver/Gold の各層の役割・利用者・データ状態を説明できる。
- [ ] ディメンショナルモデル/スタースキーマが Gold 層で作られ、正規化は主に Silver で行われることを説明できる。
- [ ] CDC が安定したサロゲートキーと下流マテビューの増分リフレッシュを可能にすることを説明できる。
- [ ] Bronze→Silver→Gold の ETL を、自動ローダー・期待値・AUTO CDC・マテビューで実装する流れを説明できる。
- [ ] CDC API に Pro/Advanced エディションが必要で、Apache Spark 宣言パイプラインでは非サポートであることを知っている。
- [ ] (発展)バイテンポラル追跡が業務時間/システム時間の 2 軸を追跡し
__SYSTEM_START_AT/__SYSTEM_END_ATを追加することを説明できる。