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

フロー(パイプライン)を作成する

パイプライン内のテーブルのフローまたはバックフィルを作成するには、 CREATE FLOWステートメントを使用します。

注記

CREATE FLOW ストリーミングテーブルをターゲットとする場合、AUTO CDC ... INTOフローとREPLACE WHEREフローの両方をサポートします。CREATE TABLE ... FLOWで作成されたマネージドテーブルはチェンジデータキャプチャをサポートしていません。マネージドテーブルに対するAUTO CDC ... INTOフローは、MANAGED_TABLE_DOES_NOT_SUPPORT_CDCで失敗します。CREATE FLOWをストリーミングテーブルに取り込む場合、その制限の対象にはなりません。CREATE TABLE ... FLOW(パイプライン)を参照してください。

構文​

CREATE FLOW flow_name [COMMENT comment] AS
{
AUTO CDC [ONCE] INTO target_table create_auto_cdc_flow_spec |
AUTO CDC [ONCE] INTO target_table create_auto_cdc_from_snapshot_spec |
INSERT [ONCE] INTO target_table BY NAME [ replace_using_spec | REPLACE WHERE condition ] query
}

create_auto_cdc_from_snapshot_spec
FROM SNAPSHOT ( snapshot_query )
[ WITH VERSION ( version_query ) ]
KEYS ( key [, ...] )
[ STORED AS { SCD TYPE 1 | SCD TYPE 2 } ]
[ TRACK HISTORY ON { col_list | * EXCEPT ( col_list ) } ]

replace_using_spec
REPLACE USING ( column_name [, ...] ) SEQUENCE BY sequence_column

パラメーター​

  • フロー名

    作成するフローの名前。

  • comment

    フローのオプションの説明。

  • AUTO CDC INTO

    create_auto_cdc_flow_specを使用してフローを定義するAUTO CDC ... INTOステートメント。AUTO CDC ... INTOステートメントまたはINSERT INTOステートメントのいずれかを含める必要があります。ソース クエリが変更データ セマンティクスを使用する場合は、 AUTO CDC ... INTO使用します。

    詳細については、 「AUTO CDC INTO (パイプライン)」を参照してください。

  • AUTO CDC … FROM スナップショット

    変更フィードを読み取るのではなく、スナップショットを比較して変更を導出する AUTO CDC ... INTO ステートメント。ソースでチェンジデータキャプチャが有効になっておらず、完全なスナップショットのみが利用可能な場合にこのフォームを使用します。ソースは、スナップショットデータを読み取る必須の FROM SNAPSHOT (snapshot_query) 句と、処理する次のスナップショットバージョンを選択するオプションの WITH VERSION (version_query) 句の 2 つの部分で指定されます。AUTO CDC FROM スナップショット の仕組みを参照してください。

    • FROM スナップショット (snapshot_query)

      必須。WITH VERSION (...) によって選択されたバージョンのスナップショットデータを読み取るクエリー。エンジンは、以前にコミットされたスナップショットと結果を比較して挿入、更新、削除を導き出し、行の識別には KEYS を、変更の保存方法の決定には STORED AS を使用して、それらをターゲットに Merge します。

      WITH VERSION (...) によって選択されたバージョンを参照するには、このクエリー内で current_snapshot_version() を呼び出します。WITH VERSION (...) が指定されていない場合、FROM SNAPSHOT (...) 内から current_snapshot_version() を呼び出すことはできません。

      WITH VERSION (...) が省略された場合、エンジンは FROM SNAPSHOT (...) を介してソースを直接読み取り、スナップショットクエリーは、ターゲットにコミット済みのデータおよびコミット済みのスナップショット状態が存在しない初期ロード中のみ実行されます。その後の更新時に、ターゲットにすでにデータが含まれているか、コミット済みのスナップショット状態が存在する場合、フローは AUTO_CDC_FROM_SNAPSHOT_NON_EMPTY_TARGET_WITHOUT_VERSION で失敗します。複数の更新にわたってスナップショットを処理するには、WITH VERSION (...) を使用します。

    • WITH VERSION (version_query)

      オプション。処理する次のスナップショット バージョンを選択するクエリー。順序付け可能な型の列を正確に1つだけ返し、行数が0または1のいずれかでなければなりません。1行を返す場合、その値はnull以外でなければなりません。列は、BIGINT などのスカラー値、またはフィールドがすべて順序付け可能な STRUCT のいずれかにすることができます。複数の列、複数の行、または null 値を返すバージョン クエリーは、INVALID_AUTO_CDC_FROM_SNAPSHOT_VERSION_QUERY でフローを失敗させます。

      1 回のパイプライン更新の間、エンジンは次のステップを繰り返します。バージョン クエリーを評価し、クエリーが 0 行を返す場合は現在の更新に対するこのフローの処理を停止し、クエリーが 1 行を返す場合は、エンジンはその値を current_snapshot_version() を通じて公開し、スナップショット クエリーを評価し、結果のスナップショットを commit し、コミットされたバージョンを last_snapshot_version() を通じて公開します。次に、エンジンはバージョン クエリーを再評価して、次のバージョンを選択します。単一のパイプライン更新では、バージョン クエリーが何も行を返さなくなるまで、バージョンが順番に処理されます。

      正常な commit の後に返されるすべてのバージョンは、以前に commit されたバージョンよりも大きくなければなりません。増加しないバージョンは、APPLY_CHANGES_FROM_SNAPSHOT_ERROR.OUT_OF_ORDER_SNAPSHOT_VERSION で更新を失敗させます。バージョン値のデータ型は、スナップショット commit 間で変更せずに維持する必要があります。データ型が変更されると、AUTO_CDC_FROM_SNAPSHOT_VERSION_SCHEMA_CHANGED で更新が失敗します。フル更新を行うと、永続化されたバージョン状態がクリアされます。

    • キー

      必須。変更検出のためにスナップショット間の行を特定するために使用される主キー列。

    • STORED AS { SCD TYPE 1 | SCD TYPE 2 }

      オプション。変更がターゲットテーブルにどのように保存されるかを指定します。defaultは SCD TYPE 1 です。

    • TRACK HISTORY ON { col_list | * EXCEPT (col_list) }

      オプション。SCD TYPE 2 でのみ適用されます。変更時に新しい履歴行をTriggerする列を指定します。明示的な列リストを指定するか、* EXCEPT (col_list) を指定して、リストされた列を除くすべての列を追跡します。

    スナップショット CDC は、WHERE または SEQUENCE BY をサポートしていません。クロススナップショットの順序付けは、WITH VERSION (...) を介して表現されます。

  • ターゲットテーブル

    更新するテーブル。これはストリーミング テーブルである必要があります。

  • INSERT INTO

    ターゲット テーブルに挿入されるテーブル クエリを定義します。ONCEオプションが指定されていない場合、クエリは ストリーミング クエリである必要があります。ストリーム キーワードを使用して、ストリーミング セマンティクスを使用してソースから読み取ります。 読み取り中に既存のレコードの変更または削除が検出されると、エラーがスローされます。静的ソースまたは追加専用のソースから読み取るのが最も安全です。変更コミットを含むデータを取り込むには、Python とskipChangeCommitsオプションを使用してエラーを処理できます。

    INSERT INTO AUTO CDC ... INTOと相互に排他的です。ソース データにチェンジデータ キャプチャ ( CDC ) 機能が含まれる場合は、 AUTO CDC ... INTO使用します。 ソースが使用していない場合はINSERT INTOを使用します。

    ストリーミング データの詳細については、 「パイプラインを使用したデータの変換」を参照してください。

  • REPLACE USING ( column_name [, ...] ) SEQUENCE BY sequence_column

備考

ベータ版

この機能はベータ版です。Databricks Runtime 18.2 以降が必要です。

フローを REPLACE USING フローとして定義します。これは、ターゲットテーブル内の指定されたキー列に一致するすべての行を置き換え、他のすべての行はそのまま残します。ソースが列でキー指定された一連の部分的なスナップショットである場合は、REPLACE USING を使用します。SEQUENCE BYは更新を順序付けし、更新が順不同で到着した場合でも、キーに対して最も高いシーケンスが優先されるようにします。

少なくとも1つのキー列と、1つのSEQUENCE BY列を指定してください。クエリーはストリーミングクエリーである必要があり、BY NAMEが必要です。REPLACE USINGは、ONCEまたはAUTO CDC ... INTOと組み合わせることはできません。

詳細については、「REPLACE USING フローによる部分的なスナップショットの置換」を参照してください。

  • REPLACE WHERE条件

    フローを REPLACE WHERE フローとして定義します。これにより、ターゲットテーブルのターゲットサブセットが再計算され、上書きされます。更新ごとに、condition に一致するターゲットテーブル内のすべての行が削除され、その同じ述語範囲に対してソースクエリーが再計算され、結果が挿入されます。condition に一致しない行はそのまま変更されずに残ります。ソースクエリーに述語を追加する必要はありません。パイプラインエンジンは、ソースからの読み込み時にそれを自動的に適用します。

    REPLACE WHERE バッチ セマンティクスを使用するため、ソース クエリーはストリーミング クエリーである必要はありません。BY NAMEは必須です。REPLACE WHEREは、ONCE、REPLACE USING、またはAUTO CDC ... INTOと組み合わせることはできません。

    For more 情報, see バッチ processing with REPLACE WHERE flows and REPLACE WHERE flows for standalone ストリーミングテーブル.

  • ONCE

    オプションで、フローをバックフィルなどの 1 回限りのフローとして定義します。ONCEを使用すると、フローは次の 2 つの方法で変化します。

    • ソースqueryまたはcreate_auto_cdc_flow_specはストリーミング テーブルではありません。
    • フローはデフォルトで1回実行されます。パイプラインが完全に更新されると、ONCE のフローが再度実行され、データが再作成されます。

    ONCE ストリーミングソースを必要とする REPLACE USING とは併用できません。

例​

SQL
-- EXAMPLE 1:
-- Create a streaming table, and add two flows that append data to it:
CREATE OR REFRESH STREAMING TABLE users;

-- first flow into target_table:
CREATE FLOW users_flow AS
INSERT INTO users BY NAME
SELECT * FROM stream(raw_data.users);

-- second flow into target_table:
CREATE FLOW backfill_users AS
INSERT ONCE INTO users BY NAME
SELECT * FROM user_backfill_table;

-- EXAMPLE 2:
-- Create a streaming table, and add a flow that applies CDC changes to it:
CREATE OR REFRESH STREAMING TABLE admins_cdc_target_table;

-- first flow into target_table:
CREATE FLOW admin_cdc_flow AS
AUTO CDC INTO admins_cdc_target_table
FROM stream(cdc_data.admins)
KEYS (userId)
APPLY AS DELETE WHEN
operation = "DELETE"
SEQUENCE BY sequenceNum
COLUMNS * EXCEPT (operation, sequenceNum)
STORED AS SCD TYPE 2;

-- EXAMPLE 3:
-- Create a streaming table, and add a REPLACE USING flow that keeps the latest
-- row for each payment_id from a stream of partial snapshots:
CREATE OR REFRESH STREAMING TABLE payments_latest;

CREATE FLOW payments_replace_flow AS
INSERT INTO payments_latest BY NAME
REPLACE USING (payment_id) SEQUENCE BY payment_date
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

-- EXAMPLE 4:
-- AUTO CDC FROM SNAPSHOT without WITH VERSION: a one-time initial load from a snapshot table.
-- To process later snapshots on each update, add WITH VERSION (see EXAMPLE 5).
CREATE STREAMING TABLE users (user_id INT, name STRING, email STRING);

CREATE FLOW users_snapshot_flow AS
AUTO CDC ONCE INTO users
FROM SNAPSHOT (SELECT * FROM catalog.schema.users_snapshot)
KEYS (user_id)
STORED AS SCD TYPE 1;

-- EXAMPLE 5:
-- AUTO CDC FROM SNAPSHOT with WITH VERSION: pick the next file, then read it as the snapshot:
CREATE STREAMING TABLE orders (order_id INT, product STRING, quantity INT, order_date DATE);

CREATE FLOW orders_cdc AS
AUTO CDC INTO orders
FROM SNAPSHOT (
SELECT order_id, product, quantity, order_date
FROM read_files('/Volumes/catalog/schema/landing/orders/', format => 'json')
WHERE _metadata.file_path = (SELECT version.path FROM current_snapshot_version())
)
WITH VERSION (
SELECT struct(modification_time, path) AS version
FROM list_files('/Volumes/catalog/schema/landing/orders/')
WHERE (
NOT EXISTS (SELECT 1 FROM last_snapshot_version())
OR struct(modification_time, path) > (SELECT version FROM last_snapshot_version())
)
ORDER BY modification_time, path
LIMIT 1
)
KEYS (order_id)
STORED AS SCD TYPE 2;

-- EXAMPLE 6:
-- One-time snapshot backfill plus a streaming CDC flow into the same target.
-- The backfill omits WITH VERSION, so it uses an implicit timestamp version. The
-- streaming flow's SEQUENCE BY column (event_ts) must be a TIMESTAMP so its type
-- matches that implicit version on the shared target.
CREATE STREAMING TABLE customers (
customer_id INT, name STRING, email STRING, address STRING, event_ts TIMESTAMP
);

CREATE FLOW customers_snapshot_backfill AS
AUTO CDC ONCE INTO customers
FROM SNAPSHOT (SELECT * FROM catalog.schema.customers_snapshot)
KEYS (customer_id)
STORED AS SCD TYPE 1;

CREATE FLOW customers_cdc AS
AUTO CDC INTO customers
FROM STREAM(customers_cdc_events)
KEYS (customer_id)
SEQUENCE BY event_ts
STORED AS SCD TYPE 1;

-- EXAMPLE 7:
-- Create a streaming table, then add a REPLACE WHERE flow that recomputes and
-- overwrites a targeted window of the target table on each update:
CREATE STREAMING TABLE payments_latest;

CREATE FLOW payments_latest AS
INSERT INTO payments_latest BY NAME
REPLACE WHERE payment_date >= date_add(current_date(), -7)
SELECT payment_id, booking_id, status, payment_date
FROM samples.wanderbricks.payments;

同じターゲット上で 1 回限りのバックフィルと継続的な CDC を組み合わせる詳細については、パイプラインを使用したヒストリカルデータのバックフィルを参照してください。

このページの見出し