使用 AUTO CDC ... INTO 语句创建使用 Lakeflow 管道更改数据捕获(CDC)功能的流。 此语句从 CDC 源读取更改,并将其应用到流式处理目标。
- 若要了解 CDC,请参阅 更改数据捕获和快照。
- 有关使用
AUTO CDC的详细信息,请参阅 AUTO CDC API:使用管道简化更改数据捕获。 - 有关更多详细信息
CREATE FLOW,请参阅 CREATE FLOW (管道)。
Syntax
CREATE OR REFRESH STREAMING TABLE table_name;
CREATE FLOW flow_name AS AUTO CDC [ONCE] INTO table_name
FROM source
KEYS (keys)
[IGNORE NULL UPDATES [ON {columnList | * EXCEPT (exceptColumnList)}]]
[APPLY AS DELETE WHEN condition]
[APPLY AS TRUNCATE WHEN condition]
SEQUENCE BY orderByColumn
[SYSTEM SEQUENCE BY systemOrderByColumn]
[COLUMNS {columnList | * EXCEPT (exceptColumnList)}]
[STORED AS {SCD TYPE 1 | SCD TYPE 2 | BITEMPORAL}]
[TRACK HISTORY ON {columnList | * EXCEPT (exceptColumnList)}]
[COLUMNS TO UPDATE columnName]
使用与其他管道查询相同的 CONSTRAINT 子句为目标定义数据质量约束。 请参阅通过管道预期管理数据质量。
默认行为对于 INSERT 和 UPDATE 事件是从源 更新插入 CDC 事件:对于符合指定键的目标表中的行进行更新,或者在目标表中不存在匹配记录时插入新行。 可以通过DELETE条件来指定APPLY AS DELETE WHEN事件的处理。
重要
您必须声明一个目标流式表,以便应用更改。 可以选择为目标表指定架构。 对于 SCD 类型 2 表,指定目标表的架构时,还必须包含与__START_AT字段具有相同数据类型的__END_ATsequence_by列。
请参阅 AUTO CDC API:使用管道简化变更数据捕获。
参数
ONCE指定
ONCE意味着这会在目标表中执行一次性插入或回填。 如果管道已刷新,则不会重新运行该管道,但在完全刷新的情况下除外。此子句是可选的。
flow_name要创建的流的名称。
source数据的源。 源必须是 流媒体 来源。 要使用流式处理语义从源中读取,请使用 STREAM 关键字。 如果读取遇到对现有记录的更改或删除,则会引发错误。 从静态源或仅限追加的源读取是最安全的。 若要引入具有更改提交的数据,可以使用 Python 和
skipChangeCommits选项来处理错误。有关流数据的详细信息,请参阅 使用管道转换数据。
KEYS用于唯一标识源数据中的行的列或列组合。 这些列中的值用于标识哪些 CDC 事件应用于目标表中的特定记录。
若要定义列的组合,请使用以逗号分隔的列列表。
此条款是必要的。
IGNORE NULL UPDATES允许导入包含目标列子集的更新。 当 CDC 事件与现有行匹配并
IGNORE NULL UPDATES指定时,具有null值的列会保留其目标中的现有值。 这也适用于具有null值的嵌套列。对于部分更新,请添加一个
ON子句来控制哪些列忽略null值:-
IGNORE NULL UPDATES ON columnList:只有列出的列在传入值时保留其现有值null。 所有其他列都应用显式null值。 -
IGNORE NULL UPDATES ON * EXCEPT (exceptColumnList):除列出的列之外的所有列在传入值时保留其现有值null。 列出的列应用显式null值。
此子句是可选的。
默认是用
null值覆盖现有列。-
APPLY AS DELETE WHEN指定何时应将 CDC 事件视为
DELETE而不是更新插入。对于 SCD 类型 2 源,为了处理无序数据,已删除的行暂时保留为基础 Delta 表中的墓碑,并在元存储中创建一个视图,用于筛选掉这些墓碑。 可以使用
pipelines.cdc.tombstoneGCThresholdInSeconds配置保留间隔。此子句是可选的。
APPLY AS TRUNCATE WHEN指定何时应将 CDC 事件视为完整表
TRUNCATE。 由于此子句会触发目标表的完全截断,因此应仅将其用于需要此功能的特定用例。APPLY AS TRUNCATE WHEN子句仅支持 SCD 类型 1。 SCD 类型 2 不支持截断操作。此子句是可选的。
SEQUENCE BY指定源数据中 CDC 事件的逻辑顺序的列名。 管道处理使用此排序来处理无序到达的更改事件。
如果需要多个列进行排序,请使用表达式:它先按第一个
STRUCT结构字段排序,然后按第二个字段进行排序(如果有平线等)。指定的列必须是可排序的数据类型。
此条款是必要的。
COLUMNS指定要包含在目标表中的列的子集。 您可以选择:
- 指定要包括的列的完整列表:
COLUMNS (userId, name, city)。 - 指定要排除的列的列表:
COLUMNS * EXCEPT (operation, sequenceNum)
此子句是可选的。
当未指定
COLUMNS子句时,默认是在目标表中包括所有列。- 指定要包括的列的完整列表:
STORED AS是将记录存储为 SCD 类型 1、SCD 类型 2 还是位。
BITEMPORAL设置为跟踪跨业务时间和系统时间的更改。 Bitemporal 需要SYSTEM SEQUENCE BY且处于 Beta 阶段。 请参阅 Bitemporal AUTO CDC。此子句是可选的。
默认值为 SCD 类型 1。
TRACK HISTORY ON指定输出列的子集,以在对这些指定列进行任何更改时生成历史记录。 您可以选择:
- 指定要跟踪的列的完整列表:
COLUMNS (userId, name, city)。 - 指定要从跟踪中排除的列的列表:
COLUMNS * EXCEPT (operation, sequenceNum)
此子句是可选的。 默认值是在发生任何更改时跟踪所有输出列的历史记录,等效于
TRACK HISTORY ON *。- 指定要跟踪的列的完整列表:
COLUMNS TO UPDATE指定要作为列名称字符串数组更新的列集,用于保存每个更改记录的源列的名称。
array<string>数组中没有的列保留其现有目标值,而列出的列是从源写入的,包括显式null值。当每个更改记录更新一组不同的列,并且需要应用显式
null值时,请使用此子句进行部分更新。不能
COLUMNS TO UPDATE与 bitemporal 表一IGNORE NULL UPDATES起使用,并且不支持它。此子句是可选的。
例子
-- Create a streaming table, then use AUTO CDC to populate it:
CREATE OR REFRESH STREAMING TABLE target;
CREATE FLOW flow
AS AUTO CDC INTO
target
FROM stream(cdc_data.users)
KEYS (userId)
APPLY AS DELETE WHEN operation = "DELETE"
SEQUENCE BY sequenceNum
COLUMNS * EXCEPT (operation, sequenceNum)
STORED AS SCD TYPE 2
TRACK HISTORY ON * EXCEPT (city);