CREATE FLOW (管道)

使用 CREATE FLOW 语句为管道流程中的表创建数据流或回填。

Syntax

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 INTO

    定义流程的语句是使用 AUTO CDC ... INTOcreate_auto_cdc_flow_spec。 必须包括 AUTO CDC ... INTO 语句或 INSERT INTO 语句。 当源查询使用更改数据语义时使用 AUTO CDC ... INTO

    有关详细信息,请参阅 AUTO CDC INTO (管道)。

  • target_table

    要更新的表。 这必须是流式处理表。

  • INSERT 转为

    定义了一个查询,该查询用于向目标表插入数据。 ONCE如果未提供该选项,则查询必须是流式处理查询。 要使用流式处理语义从源中读取,请使用 STREAM 关键字。 如果读取遇到对现有记录的更改或删除,则会引发错误。 从静态源或仅限追加的源读取是最安全的。 若要引入具有更改提交的数据,可以使用 Python 和 skipChangeCommits 选项来处理错误。

    INSERT INTOAUTO CDC ... INTO 是互斥的。 如果源数据包括更改数据捕获(CDC)功能,请使用 AUTO CDC ... INTO 。 当源文件未使用时使用INSERT INTO

    有关流数据的详细信息,请参阅 使用管道转换数据

  • 替代使用( column_name [, ...])按sequence_column顺序排列

    Important

    此功能在 Beta 版中。 需要 Databricks 运行时 18.2 及以上版本。

    定义该流程为一个 REPLACE USING 流程,替换目标表中所有匹配指定键列的行,其他行保持不动。 当你的源是按列键的部分快照系列时使用 REPLACE USINGSEQUENCE BY 更新顺序使得密钥的最高序列获胜,即使更新顺序不一致。

    至少指定一个键列,且精确指定一 SEQUENCE BY 列。 查询必须是流式查询,且 BY NAME 是必需的。 REPLACE USING 不能与 ONCEAUTO CDC ... INTO或 结合。

    更多信息请参见 “部分快照替换”使用流“

  • 一次

    可选地将流定义为一次性流,例如回填流。 通过两种方式使用 ONCE 更改流:

    • querycreate_auto_cdc_flow_spec 不是流式处理表。
    • 默认情况运行一次。 如果管道通过完全刷新进行更新,则 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);