CREATE FLOW (パイプライン)

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

パラメーター

  • flow_name

    作成するフローの名前。

  • コメント

    フローの説明 (省略可能)。

  • 自動CDCインポート

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

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

  • target_table

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

  • INSERT を に変換する

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

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

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

  • 使用を置き換えてください( column_name [、...])sequence_column順

    Important

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

    フローを REPLACE USING フローとして定義し、指定されたキーカラムに一致するターゲットテーブルのすべての行を置き換え、他の行はそのままにします。 REPLACE USING元が列ごとにキー化された部分的なスナップショットの連続であれば使います。 SEQUENCE BY 更新の順序は、キーの最高順位が勝つようにします。たとえ更新が順不同であってもです。

    少なくとも1つのキーカラムと、正確に1つの SEQUENCE BY カラムを指定します。 クエリはストリーミングクエリでなければならず、 BY NAME が必要です。 REPLACE USING ONCEAUTO CDC ... INTOと組み合わせることはできない。

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

  • ある時

    必要に応じて、バックフィルなどの 1 回限りのフローとしてフローを定義します。 ONCEを使用すると、次の 2 つの方法でフローが変更されます。

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

    ONCE REPLACE USINGではストリーミングソースが必要です。

例示

-- 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);