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

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

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

構文

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

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 (パイプライン)」を参照してください。

  • ターゲットテーブル

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

  • 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 フローによる部分的なスナップショットの置換」を参照してください。

  • 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);
このページの見出し