使用 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
要创建的流的名称。
评论
流的可选说明。
-
定义流程的语句是使用
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);