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 ... INTO 和 create_auto_cdc_flow_spec。 必须包括 AUTO 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

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

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

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

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

  • 一次

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

    • 源 query 或 create_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);