メインコンテンツまでスキップ

高度なAUTO CDCトピック

基本的な AUTO CDCAUTO CDC FROM SNAPSHOT の APIs に加えて、ターゲットテーブルで DML をランしたり、CDC ターゲットからチェンジデータフィードを読み取ったり、処理メトリクスを監視したり、部分的な更新を適用したり、バイテンポラルストレージで変更を追跡したりできます。AUTO CDC APIs の概要については、「AUTO CDC APIs: パイプラインによるチェンジデータキャプチャの簡素化」を参照してください。

ターゲットのストリーミングテーブルのデータを追加、変更、または削除します。

パイプラインがテーブルをUnity Catalogに公開する場合、insert、update、delete、merge ステートメントを含むデータ操作言語(DML) ステートメントを使用して、 AUTO CDC ... INTOステートメントによって作成されたターゲット ストリーミング テーブルを変更できます。

注記
  • ストリーミングテーブルのテーブルスキーマを変更する DML ステートメントはサポートされていません。DML ステートメントによってテーブルスキーマの進化が発生しないようにしてください。
  • ストリーミングテーブルを更新する DML ステートメントは、Databricks Runtime 13.3 LTS 以上を使用する共有Unity Catalog クラスターまたはSQL ウェアハウスでのみ実行できます。
  • ストリーミングには追加専用のデータ ソースが必要なため、処理で (DML ステートメントなどによる) 変更を伴うソース ストリーミング テーブルからのストリーミングが必要な場合は、ソース ストリーミング テーブルの読み取り時にSkipChangeCommits フラグを設定します。 skipChangeCommitsが設定されている場合、ソース テーブルのレコードを削除または変更するトランザクションは無視されます。処理にストリーミング テーブルが必要ない場合は、ターゲット テーブルとしてマテリアライズドビュー (追加のみの制限がない) を使用できます。

パイプラインは指定されたSEQUENCE BY列を使用し、ターゲットテーブルの__START_AT列と__END_AT列 (SCDタイプ 2 の場合) に適切な順序付け値を伝播するため、レコードの適切な順序を維持するために、DMLステートメントでこれらの列に有効な値が使用されていることを確認する必要があります。「AUTO CDCの仕組み」を参照してください。

ストリーミングテーブルでの DML ステートメントの使用の詳細については、 ストリーミングテーブルのデータを追加、変更、または削除するを参照してください。

次の例では、開始シーケンスを 5 としたアクティブレコードを挿入しています。

SQL
INSERT INTO my_streaming_table (id, name, __START_AT, __END_AT) VALUES (123, 'John Doe', 5, NULL);
ヒント

SCD タイプ 2 ターゲット テーブル内の__START_AT列と__END_AT列の名前を変更する必要がある場合 (たとえば、ダウンストリーム スキーマの要件に一致させるため)、ターゲット テーブルにビューを作成します。

SQL
CREATE VIEW my_employees_view AS
SELECT
*,
__START_AT AS valid_from,
__END_AT AS valid_to
FROM my_scd2_target_table;

AUTO CDCターゲットテーブルから変更データフィードを読み取ります

Databricks Runtime 15.2 以降では、他の Delta テーブルから変更データフィードを読み取るのと同じ方法で、 AUTO CDCまたはAUTO CDC FROM SNAPSHOTクエリのターゲットであるDeltaテーブルから変更データフィードを読み取ることができます。 ターゲットのストリーミング テーブルから変更データフィードを読み取るには、次のものが必要です。

  • ターゲット ストリーミングテーブルは Unity Catalogにパブリッシュする必要があります。 「パイプラインで Unity Catalog を使用する」を参照してください。
  • ターゲット ストリーミング テーブルから変更データフィードを読み取るには、 Databricks Runtime 15.2 以降を使用する必要があります。 別のパイプラインで変更データフィードを読み取るには、 Databricks Runtime 15.2 以降を使用するようにパイプラインを構成する必要があります。

LakeFlow Pipelinesで作成されたターゲットストリーミングテーブルからチェンジデータフィードを読み取る方法は、他の Delta テーブルからチェンジデータフィードを読み取るのと同じ方法です。Delta チェンジデータフィード機能の使用方法 (Python や SQL の例を含め) の詳細については、Databricks でチェンジデータフィードを使用する」を参照してください。

注記

変更データフィード レコードには、変更イベントのタイプを識別するメタデータが含まれています。 テーブル内のレコードが更新されると、関連付けられた変更レコードのメタデータには通常、 update_preimageおよびupdate_postimageイベントに設定された_change_type値が含まれます。

ただし、主キー値の変更を含む更新がターゲット ストリーミング テーブルに行われた場合、 _change_type値は異なります。 変更に主キーの更新が含まれる場合、 _change_typeメタデータ フィールドはinsertおよびdeleteイベントに設定されます。主キーの変更は、 UPDATEまたはMERGEステートメントを使用してキー フィールドの 1 つに手動で更新が行われたとき、または SCD タイプ 2 テーブルの場合は、 __start_atフィールドが以前の開始シーケンス値を反映するように変更されたときに発生する可能性があります。

AUTO CDCクエリは、SCD タイプ 1 と SCD タイプ 2 の処理で異なる主キー値を決定します。

SCD Type

Primary key

SCD type 1, and the pipelines Python interface

The primary key is the value of the keys parameter in the create_auto_cdc_flow() function. For the SQL interface the primary key is the columns defined by the KEYS clause in the AUTO CDC ... INTO statement.

SCD type 2

The primary key is the keys parameter or KEYS clause plus the return value from the coalesce(__START_AT, __END_AT) operation, where __START_AT and __END_AT are the corresponding columns from the target streaming table. This uses __START_AT when available, and __END_AT when __START_AT is null (for example, the initial record).

SCD Type

Primary key

SCD type 1, and the pipelines Python interface

The primary key is the value of the keys parameter in the create_auto_cdc_flow() function. For the SQL interface the primary key is the columns defined by the KEYS clause in the AUTO CDC ... INTO statement.

SCD type 2

The primary key is the keys parameter or KEYS clause plus the return value from the coalesce(__START_AT, __END_AT) operation, where __START_AT and __END_AT are the corresponding columns from the target streaming table. This uses __START_AT when available, and __END_AT when __START_AT is null (for example, the initial record).

パイプライン内の CDC クエリによって処理されたレコードに関するデータを取得する

注記

次のメトリクスは、 AUTO CDCクエリによってのみキャプチャされ、 AUTO CDC FROM SNAPSHOTクエリによってはキャプチャされません。

次のメトリクスは、 AUTO CDCクエリによってキャプチャされます。

  • num_upserted_rows : 更新中にデータセットにアップサートされた出力行の数。
  • num_deleted_rows : 更新中にデータセットから削除された既存の出力行の数。

非 CDC フローの出力であるnum_output_rowsメトリクスは、 AUTO CDCクエリではキャプチャされません。

部分的な更新を適用します

ソースが変更された列のみを送信する場合、AUTO CDCは、変更レコードに存在しない列(ターゲット値を変更せずに残すべき)と、nullに明示的に設定された列(ターゲット値をnullで上書きすべき)を区別する必要があります。By default, IGNORE NULL UPDATESはすべてのnullを「更新しない」マーカーとして扱うため、明示的なnullを適用できません。この曖昧さを解決するには、次の3つの方法のいずれかを選択します。

手法

使用方法

挙動

IGNORE NULL UPDATES ON columnList

小さく固定された列のセットは null の値を無視するべきであり、他のすべての列は明示的な null の値を適用します。

一覧表示されている列は、入力値がnullの場合、既存のターゲット値を保持します。その他すべての列は明示的なnull値を適用します。

IGNORE NULL UPDATES ON * EXCEPT (exceptColumnList)

ほとんどの列ではnull値を無視し、少数の列でのみ明示的なnull値を適用する必要があります。

リストされた列は、明示的なnullの値を適用します。入力値がnullの場合、他のすべての列は既存のターゲット値を保持します。

COLUMNS TO UPDATE

各変更レコードは異なる列セットを更新するか、更新可能な列セットが時間の経過とともに変化します。

ソース列は、変更レコードごとに更新する列を指定します。リストされた列は、明示的なnull値を含め、ソースから書き込まれます。リストされていない列は、既存のターゲット値を保持します。

手法

使用方法

挙動

IGNORE NULL UPDATES ON columnList

小さく固定された列のセットは null の値を無視するべきであり、他のすべての列は明示的な null の値を適用します。

一覧表示されている列は、入力値がnullの場合、既存のターゲット値を保持します。その他すべての列は明示的なnull値を適用します。

IGNORE NULL UPDATES ON * EXCEPT (exceptColumnList)

ほとんどの列ではnull値を無視し、少数の列でのみ明示的なnull値を適用する必要があります。

リストされた列は、明示的なnullの値を適用します。入力値がnullの場合、他のすべての列は既存のターゲット値を保持します。

COLUMNS TO UPDATE

各変更レコードは異なる列セットを更新するか、更新可能な列セットが時間の経過とともに変化します。

ソース列は、変更レコードごとに更新する列を指定します。リストされた列は、明示的なnull値を含め、ソースから書き込まれます。リストされていない列は、既存のターゲット値を保持します。

COLUMNS TO UPDATE IGNORE NULL UPDATESと組み合わせることはできず、バイテンポラルテーブルではサポートされていません。

経験則として、プロデューサーが各レコードでどの列が変更されたかを認識しており、その情報をソース列で伝達できる場合(複数のプロデューサーが同じソースに書き込む場合や、更新可能な列のセットが時間とともに増加する場合など)は、COLUMNS TO UPDATEを選択してください。パイプラインの所有者が更新可能な列の固定セットを事前に認識しており、パイプラインコードでそれらを制御することを好む場合は、IGNORE NULL UPDATES ONを選択してください。

次の例では、columnsToUpdateという名前のソース列を使用して、nullに明示的に設定された列を含む、各変更レコードが更新する列を制御します:

Python
from pyspark import pipelines as dp

dp.create_streaming_table("target")

dp.create_auto_cdc_flow(
target = "target",
source = "cdc_source",
keys = ["id"],
sequence_by = "sequenceNum",
stored_as_scd_type = 1,
columns_to_update = "columnsToUpdate"
)

完全なパラメーターリファレンスについては、「 AUTO CDC INTO (パイプライン) 」および「 create_auto_cdc_flow 」を参照してください。

バイテンポラルAUTO CDC

備考

ベータ版

バイテンポラル AUTO CDC はベータ版です。

SCD タイプ 1 とタイプ 2 は単一時間軸です。これらは単一の時間ディメンションにわたる変更を追跡します。バイテンポラルは、SCD タイプ 2 の履歴を拡張し、2 つの時間ディメンションにわたる変更を追跡し、2 つの視点を区別します:

  • **ビジネス時間**:イベントが実際に発生したとき。
  • システム時間 :システムがイベントを記録または取り込んだ時刻。

SCDタイプ2と同様に、バイテンポラルはレコードの完全な履歴を保持します。これにより2つ目のタイムラインが追加され、過去の任意の時点においてデータが何を示していたか、またシステムが何を認識していたかの両方を再構築できます。

たとえば、ヘッジファンドはソースシステムから株式データを取り込みます。Acme Corpの株価は1月1日に変更されますが、ファンドは1月5日までその更新を取り込みません。バイテンポラルAUTO CDCを使用すると、ファンドは2つの異なる質問に答えることができます。Acme Corpの実際の株価が1月1日(ビジネス時間)にいくらだったか、およびファンドが1月3日(システム時間)に取引の決定を下したときにシステムが信じていた価格です。これらのタイムラインを区別できることは、監査、規制報告、財務上の意思決定に役立ちます。

バイテンポラル処理を有効にするには、STORED AS BITEMPORAL(SQL)またはstored_as_scd_type="bitemporal"(Python)を設定し、ビジネス時間列にはSEQUENCE BYを、システム時間列にはSYSTEM SEQUENCE BYを使用します。ターゲットテーブルには__SYSTEM_START_ATおよび__SYSTEM_END_AT列が追加され、SCD Type 2の__START_ATおよび__END_AT列と並んで配置されます。構文の詳細については、「AUTO CDC INTO(パイプライン)」または「create_auto_cdc_flow」を参照してください。

バイテンポラルAUTO CDCの例

次の例では、少量の合成CDCイベントからバイテンポラルターゲットテーブルを作成します。bt列にはビジネスタイムがあり、st列にはシステムタイムがあります。

Python
from pyspark import pipelines as dp

# Source: synthetic CDC events
dp.create_streaming_table(name="cdc_source")

@dp.append_flow(target="cdc_source", once=True)
def load_cdc_source():
return spark.createDataFrame(
[
(1, "x10", "y10", 10, 100),
(1, "x20", "y20", 20, 200)
],
schema="id INT, x STRING, y STRING, bt INT, st INT",
)

# Target: bitemporal table
dp.create_streaming_table(name="target_bitemporal")

dp.create_auto_cdc_flow(
target = "target_bitemporal",
source = "cdc_source",
keys = ["id"],
sequence_by = "bt",
system_sequence_by = "st",
stored_as_scd_type = "bitemporal"
)

次の変更シーケンスは、バイテンポラルテーブルが単一の企業に対する挿入、更新、順序外更新、および削除をどのように記録するかを示しています。シーケンシング列は__START_ATおよび__END_AT(ビジネス時間)列を生成し、システムシーケンシング列は__SYSTEM_START_ATおよび__SYSTEM_END_AT(システム時間)列を生成します。

説明

__START_AT

この行が有効になった業務時間。

__END_AT

この行の有効期限が終了する業務時間です。無期限に有効な場合はnullです。

__SYSTEM_START_AT

この行のデータとビジネス時間間隔が真であると認識されているシステム時間です。

__SYSTEM_END_AT

この行のデータとビジネスの期間が無効であることが認識されているシステム時刻です。null が無期限に真であることがわかっている場合。

説明

__START_AT

この行が有効になった業務時間。

__END_AT

この行の有効期限が終了する業務時間です。無期限に有効な場合はnullです。

__SYSTEM_START_AT

この行のデータとビジネス時間間隔が真であると認識されているシステム時間です。

__SYSTEM_END_AT

この行のデータとビジネスの期間が無効であることが認識されているシステム時刻です。null が無期限に真であることがわかっている場合。

システムは、両方のタイムラインにわたって、どの順序で到着するイベントも処理します。イベントが既に処理済みのイベントよりも前のビジネス時間またはシステム時間で到着した場合、システムは末尾に追加するだけではなく、影響を受ける履歴を修正します。

変更1:挿入

会社Aは2025年7月18日10:01:00(ビジネス時間)に追加されましたが、10:05:00(システム時間)まで取り込まれません。

入力:

CompanyId

データポイント

シーケンシング

システムのシーケンス

オペレーション

A

XFv1

2025年7月18日 10時01分00秒

2025/07/18 10:05:00

INSERT

CompanyId

データポイント

シーケンシング

システムのシーケンス

オペレーション

A

XFv1

2025年7月18日 10時01分00秒

2025/07/18 10:05:00

INSERT

出力:

CompanyId

データポイント

__START_AT

__END_AT

__SYSTEM_START_AT

__SYSTEM_END_AT

A

XFv1

2025年7月18日 10時01分00秒

NULL

2025/07/18 10:05:00

NULL

CompanyId

データポイント

__START_AT

__END_AT

__SYSTEM_START_AT

__SYSTEM_END_AT

A

XFv1

2025年7月18日 10時01分00秒

NULL

2025/07/18 10:05:00

NULL

XFv1 は 10:01:00 から有効で、終了時期は未定です。システムはシステム時刻 10:05:00 にこの事実を認識し、終了時期は未定です。

変更 2: 更新

会社Aは2025年7月18日 12時15分43秒(ビジネス時間)に更新され、システムは12時20分00秒(システム時間)にイベントを処理します。システムは、更新が認識される前にシステムが信じていたことと、更新が取り込まれた後の修正されたビジネス履歴の両方を保持します。

入力:

CompanyId

データポイント

シーケンシング

システムのシーケンス

オペレーション

A

XFv2

2025/07/18 12:15:43

2025年7月18日 12時20分00秒

UPDATE

CompanyId

データポイント

シーケンシング

システムのシーケンス

オペレーション

A

XFv2

2025/07/18 12:15:43

2025年7月18日 12時20分00秒

UPDATE

出力:

CompanyId

データポイント

__START_AT

__END_AT

__SYSTEM_START_AT

__SYSTEM_END_AT

A

XFv1

2025年7月18日 10時01分00秒

NULL

2025/07/18 10:05:00

2025年7月18日 12時20分00秒

A

XFv1

2025年7月18日 10時01分00秒

2025/07/18 12:15:43

2025年7月18日 12時20分00秒

NULL

A

XFv2

2025/07/18 12:15:43

NULL

2025年7月18日 12時20分00秒

NULL

CompanyId

データポイント

__START_AT

__END_AT

__SYSTEM_START_AT

__SYSTEM_END_AT

A

XFv1

2025年7月18日 10時01分00秒

NULL

2025/07/18 10:05:00

2025年7月18日 12時20分00秒

A

XFv1

2025年7月18日 10時01分00秒

2025/07/18 12:15:43

2025年7月18日 12時20分00秒

NULL

A

XFv2

2025/07/18 12:15:43

NULL

2025年7月18日 12時20分00秒

NULL

XFv1 は、10:01:00 から既知の終了なしで有効であると信じられており、システムは 10:05:00 から 12:20:00 までその信念を保持していました。XFv1 は、12:15:43 までのみ有効であることが判明しました。これは、システム時刻 12:20:00 から既知の終了なしで有効な、修正された履歴です。XFv2 は、12:15:43 から既知の終了なしで有効であり、システム時刻 12:20:00 に学習されました。

変更3:順不同の更新

順不同の更新が到着し、Company A が実際に 2025/07/18 12:05:00 (ビジネス時間) に更新されたことを示していますが、12:25:00 (システム時間) まで取り込まれません。更新がシステム時間では遅れて到着しても、以前のビジネス時間を伴う場合、システムは履歴ビジネス時間を修正し、順不同の更新前にシステムが認識していたものと、修正された履歴の両方を保持します。

入力:

CompanyId

データポイント

シーケンシング

システムのシーケンス

オペレーション

A

XFv3

2025/07/18 12:05:00

2025年7月18日 12時25分00秒

UPDATE

CompanyId

データポイント

シーケンシング

システムのシーケンス

オペレーション

A

XFv3

2025/07/18 12:05:00

2025年7月18日 12時25分00秒

UPDATE

出力:

CompanyId

データポイント

__START_AT

__END_AT

__SYSTEM_START_AT

__SYSTEM_END_AT

A

XFv1

2025年7月18日 10時01分00秒

NULL

2025/07/18 10:05:00

2025年7月18日 12時20分00秒

A

XFv1

2025年7月18日 10時01分00秒

2025/07/18 12:15:43

2025年7月18日 12時20分00秒

2025年7月18日 12時25分00秒

A

XFv1

2025年7月18日 10時01分00秒

2025/07/18 12:05:00

2025年7月18日 12時25分00秒

NULL

A

XFv3

2025/07/18 12:05:00

2025/07/18 12:15:43

2025年7月18日 12時25分00秒

NULL

A

XFv2

2025/07/18 12:15:43

NULL

2025年7月18日 12時20分00秒

NULL

CompanyId

データポイント

__START_AT

__END_AT

__SYSTEM_START_AT

__SYSTEM_END_AT

A

XFv1

2025年7月18日 10時01分00秒

NULL

2025/07/18 10:05:00

2025年7月18日 12時20分00秒

A

XFv1

2025年7月18日 10時01分00秒

2025/07/18 12:15:43

2025年7月18日 12時20分00秒

2025年7月18日 12時25分00秒

A

XFv1

2025年7月18日 10時01分00秒

2025/07/18 12:05:00

2025年7月18日 12時25分00秒

NULL

A

XFv3

2025/07/18 12:05:00

2025/07/18 12:15:43

2025年7月18日 12時25分00秒

NULL

A

XFv2

2025/07/18 12:15:43

NULL

2025年7月18日 12時20分00秒

NULL

XFv1は10:01:00から12:15:43まで有効であると信じられていましたが、その認識は現在、システム時間で12:25:00まで有効です。新しいアップデートにより、XFv1のビジネス上の有効期限が12:05:00に修正され、この修正された履歴はシステム時間12:25:00から有効です。XFv3は現在、12:05:00から12:15:43まで有効であることが知られており、この認識はシステム時間では12:25:00から有効で、終了は不明です。

変更 4:削除

会社 A は 2025/07/18 12:30:00 に削除され、システムは 12:30:00 にそのイベントを取り込みます。削除操作はエンティティの事業上の存在の終了を表すため、システムは置換行を作成しません。XFv2 は 2 つの行に表示され、会社が消滅したときと、システムが削除を認識したときの両方の完全な監査証跡を保持します。

入力:

CompanyId

データポイント

シーケンシング

システムのシーケンス

オペレーション

A

XFv2

2025年7月18日 12時30分00秒

2025年7月18日 12時30分00秒

DELETE

CompanyId

データポイント

シーケンシング

システムのシーケンス

オペレーション

A

XFv2

2025年7月18日 12時30分00秒

2025年7月18日 12時30分00秒

DELETE

出力:

CompanyId

データポイント

__START_AT

__END_AT

__SYSTEM_START_AT

__SYSTEM_END_AT

A

XFv1

2025年7月18日 10時01分00秒

NULL

2025/07/18 10:05:00

2025年7月18日 12時20分00秒

A

XFv1

2025年7月18日 10時01分00秒

2025/07/18 12:15:43

2025年7月18日 12時20分00秒

2025年7月18日 12時25分00秒

A

XFv1

2025年7月18日 10時01分00秒

2025/07/18 12:05:00

2025年7月18日 12時25分00秒

NULL

A

XFv3

2025/07/18 12:05:00

2025/07/18 12:15:43

2025年7月18日 12時25分00秒

NULL

A

XFv2

2025/07/18 12:15:43

NULL

2025年7月18日 12時20分00秒

2025年7月18日 12時30分00秒

A

XFv2

2025/07/18 12:15:43

2025年7月18日 12時30分00秒

2025年7月18日 12時30分00秒

NULL

CompanyId

データポイント

__START_AT

__END_AT

__SYSTEM_START_AT

__SYSTEM_END_AT

A

XFv1

2025年7月18日 10時01分00秒

NULL

2025/07/18 10:05:00

2025年7月18日 12時20分00秒

A

XFv1

2025年7月18日 10時01分00秒

2025/07/18 12:15:43

2025年7月18日 12時20分00秒

2025年7月18日 12時25分00秒

A

XFv1

2025年7月18日 10時01分00秒

2025/07/18 12:05:00

2025年7月18日 12時25分00秒

NULL

A

XFv3

2025/07/18 12:05:00

2025/07/18 12:15:43

2025年7月18日 12時25分00秒

NULL

A

XFv2

2025/07/18 12:15:43

NULL

2025年7月18日 12時20分00秒

2025年7月18日 12時30分00秒

A

XFv2

2025/07/18 12:15:43

2025年7月18日 12時30分00秒

2025年7月18日 12時30分00秒

NULL

XFv2は12:15:43から有効で、既知の終了時刻はなく、システムはその状態を12:20:00から12:30:00まで認識していました。削除が取り込まれると、XFv2は12:30:00までのみ有効であることが認識され、システム時刻12:30:00から訂正された履歴が有効になります。

パイプラインの CDC 処理に使用されるデータ オブジェクトは何ですか?

Hive metastore でターゲットテーブルを宣言すると、2 つのデータ構造が作成されます。

  • ターゲットテーブルに割り当てられた名前を使用するビュー。
  • CDC 処理を管理するためにパイプラインによって使用される内部バッキング テーブル。このテーブルの名前は、ターゲット テーブル名の前に__apply_changes_storage_を付加して付けられます。

たとえば、 dp_cdc_targetという名前のターゲット テーブルを宣言すると、メタストアにdp_cdc_targetという名前のビューと__apply_changes_storage_dp_cdc_targetという名前のテーブルが表示されます。処理されたデータにアクセスするには、ビューをクエリします。バッキングテーブルを直接変更しないでください。

注記

これらのデータ構造はAUTO CDC処理にのみ適用され、 AUTO CDC FROM SNAPSHOT処理には適用されません。これらは、 Unity Catalogではなく、 Hive metastoreにのみ適用されます。