高级 AUTO CDC 主题内容

除了基本的 AUTO CDCAUTO CDC FROM SNAPSHOT API 之外,您还可以对目标表执行 DML 操作、从 CDC 目标读取变更数据馈送、监控处理指标、应用部分更新,以及使用双时态存储跟踪更改。 有关 AUTO CDC API 接口的简介,请参阅 AUTO CDC API 接口:使用管道简化变更数据捕获

在目标流式处理表中添加、更改或删除数据

如果管道将表发布到 Unity 目录,则可以使用 数据作语言 (DML)语句(包括插入、更新、删除和合并语句)修改由 AUTO CDC ... INTO 语句创建的目标流式处理表。

注释

  • 不支持用于修改流式处理表的表架构的 DML 语句。 确保 DML 语句不会尝试修改表架构。
  • 只能在使用 Databricks Runtime 13.3 LTS 及更高版本的 Unity Catalog 共享群集或 SQL 仓库中运行用于更新流式表的 DML 语句。
  • 由于流式传输要求仅追加数据源,因此如果你的处理需要从包含更改的源流式传输表进行流式传输(例如,通过 DML 语句),请在读取源流式传输表时设置 skipChangeCommits 标志。 设置 skipChangeCommits 后,将会忽略删除或修改源表上记录的事务。 如果你的处理不需要某个流式处理表,则可使用具体化视图(没有“仅追加”限制)作为目标表。

由于该管道使用指定的 SEQUENCE BY 列,并将适当的排序值传播到目标表的 __START_AT 列和 __END_AT 列中(对于 SCD 类型 2),因此您必须确保 DML 语句为这些列使用有效值,以保持记录的正确顺序。 请参阅 AUTO CDC 的工作原理

有关对流式处理表使用 DML 语句的详细信息,请参阅 在流式处理表中添加、更改或删除数据

以下示例插入一个活动记录,其起始序列为 5:

INSERT INTO my_streaming_table (id, name, __START_AT, __END_AT) VALUES (123, 'John Doe', 5, NULL);

小窍门

如果需要重命名目标表中 SCD 类型 2 的 __START_AT__END_AT 列(例如,若要匹配下游架构要求),请基于目标表创建视图:

CREATE VIEW my_employees_view AS
SELECT
  *,
  __START_AT AS valid_from,
  __END_AT AS valid_to
FROM my_scd2_target_table;

从 AUTO CDC 目标表读取变更数据流

在 Databricks Runtime 15.2 及以上版本中,你可以从作为 AUTO CDCAUTO CDC FROM SNAPSHOT 查询目标的流式处理表中读取更改数据馈送,方式与从其他 Delta 表读取更改数据馈送相同。 从目标流式处理表读取更改数据馈送需要满足以下条件:

  • 目标流式处理表必须发布到 Unity Catalog。 请参阅 将 Unity Catalog 与管道结合使用
  • 为了从目标流式处理表中读取变更数据流,您必须使用 Databricks Runtime 15.2 或更高版本。 若要读取其他管道中的更改数据馈送,必须将管道配置为使用 Databricks Runtime 15.2 或更高版本。

读取在 Lakeflow 管道中创建的目标流式表的变更数据馈送,其方式与读取其他 Delta 表的变更数据馈送相同。 要详细了解如何使用 Delta 更改数据馈送功能(包括 Python 和 SQL 示例),请参阅在 Azure Databricks 上使用更改数据馈送

注释

更改数据馈送记录包括标识更改事件的类型的 元数据 。 在表中更新记录时,关联更改记录的元数据通常包括 _change_type 值,以及设置为 update_preimageupdate_postimage 的事件。

但是,如果对目标流表进行更新并包含更改主键值,则这些 _change_type 值会有所不同。 更改包括对主键的更新时,元数据 _change_type 字段将设置为 insertdelete 事件。 在对具有 UPDATEMERGE 语句的某个键字段进行手动更新时,或对于 SCD 类型 2 表,当 __start_at 字段更改以反映较早的起始序列值时,主键可能会发生更改。

查询 AUTO CDC 用于确定主键值,而在确定这些主键值的方法上,SCD 类型 1 和 SCD 类型 2 的处理是不同的。

SCD 类型 主密钥
SCD 类型 1 和管道 Python 接口 主键是keys函数中create_auto_cdc_flow()参数的值。 对于 SQL 接口,主键是由KEYS语句中的AUTO CDC ... INTO子句定义的列。
SCD 类型 2 主键由 keys 参数或 KEYS 子句和操作的 coalesce(__START_AT, __END_AT) 返回值组成,其中 __START_AT__END_AT 是目标流式处理表中的相应列。 可用时使用 __START_AT;当 __END_AT 为 null 时,使用 __START_AT(例如,初始记录)。

获取有关管道中 CDC 查询处理的记录的数据

注释

以下指标仅由 AUTO CDC 查询捕获,而不是由 AUTO CDC FROM SNAPSHOT 查询捕获。

以下指标由 AUTO CDC 查询捕获:

  • num_upserted_rows:更新期间插入数据集的输出行数。
  • num_deleted_rows:更新期间从数据集中删除的现有输出行数。

对于 num_output_rows 查询,不会捕获 AUTO CDC 非 CDC 流的指标(输出)。

应用部分更新

当源仅发送发生更改的列时,AUTO CDC 必须能够区分两种情况:一种是某列在变更记录中缺失,此时应保持目标值不变;另一种是某列被显式设置为 null,此时应使用 null 覆盖目标值。 默认情况下, IGNORE NULL UPDATES 将每个 null 标记视为“不更新”标记,因此它不能应用显式 null标记。 若要解决此歧义,请选择以下三种方法之一:

方法 何时使用 Behavior
IGNORE NULL UPDATES ON columnList 一组较小的固定列应忽略 null 值,而所有其他列都应用显式 null 值。 当传入值为 null 时,列出的列将保留其当前目标值。 所有其他列都应用显式 null 值。
IGNORE NULL UPDATES ON * EXCEPT (exceptColumnList) 大多数列应忽略 null 值,只有少数列应应用显式 null 值。 列出的列应用了显式的 null 值。 当传入值为 null 时,所有其他列保留其现有目标值。
COLUMNS TO UPDATE 每个更改记录都会更新一组不同的列,或者一组可更新的列会随时间而更改。 源列指定针对每条更改记录要更新的列。 列出的列是从源写入的,包括显式 null 值。 未列出的列保留其现有目标值。

COLUMNS TO UPDATE 不能与 IGNORE NULL UPDATES 结合使用,并且不支持双时态表。

根据经验法则,当生产者知道每条记录中哪些列发生了更改,并且可以在源列中携带这些信息时,应选择 COLUMNS TO UPDATE;例如,当多个生产者写入同一来源,或者可更新列的集合会随着时间推移而扩大时。 当管道所有者预先知道固定的一组可更新列,并且希望在管道代码中控制这些列时,请选择 IGNORE NULL UPDATES ON

以下示例使用一个名为 columnsToUpdate 的源列来控制每条更改记录更新哪些列,包括显式设置为 null 的列:

Python

from pyspark import pipelines as dp

dp.create_streaming_table("target")

dp.create_auto_cdc_flow(
  target = "target",
  source = "cdc_source",
  keys = ["id"],
  sequence_by = "sequenceNum",
  stored_as_scd_type = 1,
  columns_to_update = "columnsToUpdate"
)

SQL

CREATE OR REFRESH STREAMING TABLE target;

CREATE FLOW apply_cdc AS AUTO CDC INTO
  target
FROM
  stream(cdc_source)
KEYS
  (id)
SEQUENCE BY
  sequenceNum
STORED AS
  SCD TYPE 1
COLUMNS TO UPDATE
  columnsToUpdate;

有关完整的参数参考,请参阅 AUTO CDC INTO (pipelines)create_auto_cdc_flow

哪些数据对象用于管道中的 CDC 处理?

在 Hive 元存储中声明目标表时,会创建两个数据结构:

  • 一个使用分配给目标表的名称的视图。
  • 管道用于管理 CDC 处理的内部后盾表。 此表通过追加 __apply_changes_storage_ 到目标表名称进行命名。

例如,如果声明名为 dp_cdc_target的目标表,则会看到一个命名 dp_cdc_target 视图和元存储中命名 __apply_changes_storage_dp_cdc_target 的表。 查询视图以访问已处理的数据。 请勿直接修改后盾表。

注释

这些数据结构仅适用于 AUTO CDC 处理,不适用于 AUTO CDC FROM SNAPSHOT 处理。 它们也仅适用于 Hive Metastore,而不适用于 Unity Catalog。