Lakeflow 管道借助 AUTO CDC 和 AUTO CDC FROM SNAPSHOT API 简化了变更数据捕获(CDC)。 这些 API 通过 CDC 源或数据库快照自动化计算渐变维度 (SCD) 类型 1 和类型 2 的复杂过程。 该 AUTO CDC API 还支持双时间跟踪(Beta),可记录跨两个时间维度的变更。 若要详细了解 SCD 类型 1 和类型 2,请参阅更改数据捕获和快照。 有关双时态跟踪的详细信息,请参阅 Bitemporal AUTO CDC。
注释
AUTO CDC API 将替换 APPLY CHANGES API,并具有相同的语法。
APPLY CHANGES API 仍然可用,但 Databricks 建议在其位置使用 AUTO CDC API。
使用的 API 取决于更改数据的源:
-
AUTO CDC:当源数据库启用了 CDC 数据流时,请使用此选项。AUTO CDC处理更改数据源 (CDF) 的更改。 管道 SQL 和 Python 接口都支持它。 -
AUTO CDC FROM SNAPSHOT:如果未在源数据库上启用 CDC,并且只有快照可用,请使用此选项。 此 API 比较快照以确定更改,然后处理这些更改。 它仅在 Python 接口中受支持。
这两个 API 都支持使用 SCD 类型 1 和类型 2 更新表:
- 使用 SCD 类型 1 直接更新记录。 不会保留已更新记录的历史记录。
- 使用 SCD 类型 2 保留有关所有更新或者对指定列集的更新的记录的历史记录。
仅对于 AUTO CDC,您还可以使用双时态存储,它将 SCD 2 型历史记录扩展为跨两个时间维度跟踪更改:业务时间和系统时间。 Bitemporal 目前为 测试版。 请参阅 Bitemporal AUTO CDC。
AUTO CDC 还支持部分更新,其中更改记录仅更新列的子集。 请参阅 “应用部分更新”。
AUTO CDC Apache Spark 声明性管道不支持这些 API。
有关语法和其他引用,请参阅 AUTO CDC INTO(管道)、create_auto_cdc_flow 和 create_auto_cdc_from_snapshot_flow。
注释
本页介绍如何根据源数据的变化来更新管道中的表。 若要了解如何记录和查询 Delta 表的行级更改信息,请参阅在Azure Databricks上使用更改数据馈送。
要求
若要使用 CDC API,必须将您的管道配置为使用 无服务器 Lakeflow 管道 或 Lakeflow 管道 Pro 或 Advanced版本。
AUTO CDC 的工作原理
若要执行 AUTO CDC CDC 处理,请首先创建流式处理表,然后分别在 SQL 中使用 AUTO CDC ... INTO 语句或在 Python 中使用 create_auto_cdc_flow() 函数来指定变更跟踪源、关键字和排序。 有关序列化和 SCD 逻辑的工作原理的说明,请参阅 更改数据捕获和快照。 请参阅 AUTO CDC 示例。
若要从具有更改源的源进行初始数据注入,请使用 AUTO CDC 和 once 流,然后继续处理更改源。 请参阅 使用 AUTO CDC 复制外部 RDBMS 表。
有关语法详细信息,请参阅 AUTO CDC INTO(管道)或 create_auto_cdc_flow。
AUTO CDC FROM SNAPSHOT 的运行机制
AUTO CDC FROM SNAPSHOT 通过比较顺序快照来确定源数据的更改。 它仅在 Python 管道接口中受支持。 可以直接从 Delta 表、云存储文件或 JDBC 读取快照。
若要使用 AUTO CDC FROM SNAPSHOT 执行 CDC 处理,请创建流式处理表,然后使用 create_auto_cdc_from_snapshot_flow() 函数指定快照、键和其他参数。 有关两种引入模式以及何时使用每个模式的详细信息,请参阅 快照处理模式。 请参阅 AUTO CDC FROM SNAPSHOT 示例。
有关语法详细信息,请参阅 create_auto_cdc_from_snapshot_flow。
使用多个列进行排序
若要按多个列(例如时间戳和 ID 以解决冲突)进行排序,请使用STRUCT来组合它们。 API先按第一个字段进行排序,如果出现相同情况,则考虑第二个字段,以此类推。
SQL
SEQUENCE BY STRUCT(timestamp_col, id_col)
Python
sequence_by = struct("timestamp_col", "id_col")
AUTO CDC 示例
以下示例演示使用更改数据馈送源处理 SCD 类型 1 和类型 2。 示例数据创建新的用户记录、删除用户记录和更新用户记录。 在 SCD 类型 1 示例中,最后的 UPDATE 个操作由于延迟到达而被从目标表中删除,说明了无序事件处理。
下面是这些示例中使用的输入记录。 通过在 “创建示例数据 ”部分中运行查询来创建此数据。
| userId | 姓名 | city | 操作 | 序列号 |
|---|---|---|---|---|
| 124 | 劳尔 | 瓦哈卡州 | INSERT | 1 |
| 123 | 伊莎贝尔 | 蒙特雷 | INSERT | 1 |
| 125 | 梅赛德斯 | 提 华纳 | INSERT | 2 |
| 126 | 百合 | 坎昆 | INSERT | 2 |
| 123 | null | null | DELETE | 6 |
| 125 | 梅赛德斯 | 瓜达拉哈拉 | UPDATE | 6 |
| 125 | 梅赛德斯 | Mexicali | UPDATE | 5 |
| 123 | 伊莎贝尔 | 奇瓦瓦州 | UPDATE | 5 |
如果在示例数据生成查询中取消注释最后一行,它将插入以下记录,指定在 sequenceNum=3 位置进行表截断(清空表):
| userId | 姓名 | city | 操作 | 序列号 |
|---|---|---|---|---|
| null | null | null | 截断 | 3 |
注释
以下所有示例都包含指定DELETE和TRUNCATE操作的选项,但每个操作都是可选的。
创建示例数据
运行以下语句以创建示例数据集。 此代码不应作为管道定义的一部分运行。 从管道的浏览文件夹运行它,而不是转换文件夹。
CREATE SCHEMA IF NOT EXISTS main.cdc_tutorial;
CREATE TABLE main.cdc_tutorial.users_cdf
AS SELECT
col1 AS userId,
col2 AS name,
col3 AS city,
col4 AS operation,
col5 AS sequenceNum
FROM (
VALUES
-- Initial load.
(124, "Raul", "Oaxaca", "INSERT", 1),
(123, "Isabel", "Monterrey", "INSERT", 1),
-- New users.
(125, "Mercedes", "Tijuana", "INSERT", 2),
(126, "Lily", "Cancun", "INSERT", 2),
-- Isabel is removed from the system and Mercedes moved to Guadalajara.
(123, null, null, "DELETE", 6),
(125, "Mercedes", "Guadalajara", "UPDATE", 6),
-- This batch of updates arrived out of order. The batch at sequenceNum 6 is the final state.
(125, "Mercedes", "Mexicali", "UPDATE", 5),
(123, "Isabel", "Chihuahua", "UPDATE", 5)
-- Uncomment to test TRUNCATE.
-- ,(null, null, null, "TRUNCATE", 3)
);
处理 SCD 类型 1 更新
SCD 类型 1 仅保留每个记录的最新版本。 以下示例从上面创建的更改数据馈送中读取数据,并将更改应用于流处理表目标。 你需要一个流水线来运行这些代码。 请参阅 使用 Lakeflow 管道编辑器开发和调试 ETL 管道。
Python
from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr
@dp.view
def users():
return spark.readStream.table("main.cdc_tutorial.users_cdf")
dp.create_streaming_table("users_current")
dp.create_auto_cdc_flow(
target = "users_current",
source = "users",
keys = ["userId"],
sequence_by = col("sequenceNum"),
apply_as_deletes = expr("operation = 'DELETE'"),
apply_as_truncates = expr("operation = 'TRUNCATE'"),
except_column_list = ["operation", "sequenceNum"],
stored_as_scd_type = 1
)
SQL
CREATE OR REFRESH STREAMING TABLE users_current;
CREATE FLOW apply_cdc AS AUTO CDC INTO
users_current
FROM
stream(main.cdc_tutorial.users_cdf)
KEYS
(userId)
APPLY AS DELETE WHEN
operation = "DELETE"
APPLY AS TRUNCATE WHEN
operation = "TRUNCATE"
SEQUENCE BY
sequenceNum
COLUMNS * EXCEPT
(operation, sequenceNum)
STORED AS
SCD TYPE 1;
运行 SCD 类型 1 示例后,目标表包含以下记录:
| userId | 姓名 | city |
|---|---|---|
| 124 | 劳尔 | 瓦哈卡州 |
| 125 | 梅赛德斯 | 瓜达拉哈拉 |
| 126 | 百合 | 坎昆 |
用户 123 (Isabel) 已被删除,并且未显示。 用户 125 (梅赛德斯) 仅显示最新的城市 (瓜达拉哈拉), 因为 SCD 类型 1 覆盖以前的值。 之前在 UPDATE 的 sequenceNum=5 被删除,因为在 sequenceNum=6 的更新已经到来。
在运行未注释的 TRUNCATE 记录的示例后,表在 sequenceNum=3 处被清除。 这意味着记录 124 和 126 不在表中,最终目标表仅包含以下记录:
| userId | 姓名 | city |
|---|---|---|
| 125 | 梅赛德斯 | 瓜达拉哈拉 |
处理 SCD 类型 2 更新
SCD 类型 2 通过为每个记录__START_AT__END_AT版本创建新行以及指示每个版本处于活动状态的列来保留更改的完整历史记录。
Python
from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr
@dp.view
def users():
return spark.readStream.table("main.cdc_tutorial.users_cdf")
dp.create_streaming_table("users_history")
dp.create_auto_cdc_flow(
target = "users_history",
source = "users",
keys = ["userId"],
sequence_by = col("sequenceNum"),
apply_as_deletes = expr("operation = 'DELETE'"),
except_column_list = ["operation", "sequenceNum"],
stored_as_scd_type = "2"
)
SQL
CREATE OR REFRESH STREAMING TABLE users_history;
CREATE FLOW apply_cdc AS AUTO CDC INTO
users_history
FROM
stream(main.cdc_tutorial.users_cdf)
KEYS
(userId)
APPLY AS DELETE WHEN
operation = "DELETE"
SEQUENCE BY
sequenceNum
COLUMNS * EXCEPT
(operation, sequenceNum)
STORED AS
SCD TYPE 2;
运行 SCD 类型 2 示例后,目标表包含以下记录:
| userId | 姓名 | city | __START_AT | __END_AT |
|---|---|---|---|---|
| 123 | 伊莎贝尔 | 蒙特雷 | 1 | 5 |
| 123 | 伊莎贝尔 | 奇瓦瓦州 | 5 | 6 |
| 124 | 劳尔 | 瓦哈卡州 | 1 | null |
| 125 | 梅赛德斯 | 提 华纳 | 2 | 5 |
| 125 | 梅赛德斯 | Mexicali | 5 | 6 |
| 125 | 梅赛德斯 | 瓜达拉哈拉 | 6 | null |
| 126 | 百合 | 坎昆 | 2 | null |
该表保留完整的历史记录。 用户 123 有两个版本(删除时结束于序列 6)。 用户 125 有三个显示城市更改的版本。 带 __END_AT = null 的记录当前处于活动状态。
使用 SCD 类型 2 跟踪列子集
默认情况下,SCD 类型 2 会在任何列值发生更改时创建新版本。 可以指定要跟踪的列的子集,以便对其他列所做的更改将更新当前版本,而不是生成新的历史记录。
以下示例从历史记录跟踪中排除 city 列:
Python
from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr
@dp.view
def users():
return spark.readStream.table("main.cdc_tutorial.users_cdf")
dp.create_streaming_table("users_history")
dp.create_auto_cdc_flow(
target = "users_history",
source = "users",
keys = ["userId"],
sequence_by = col("sequenceNum"),
apply_as_deletes = expr("operation = 'DELETE'"),
except_column_list = ["operation", "sequenceNum"],
stored_as_scd_type = "2",
track_history_except_column_list = ["city"]
)
SQL
CREATE OR REFRESH STREAMING TABLE users_history;
CREATE FLOW apply_cdc AS AUTO CDC INTO
users_history
FROM
stream(main.cdc_tutorial.users_cdf)
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)
由于 city 未跟踪更改,城市更新将覆盖当前行,而不是创建新版本。 目标表包含以下记录:
| userId | 姓名 | city | __START_AT | __END_AT |
|---|---|---|---|---|
| 123 | 伊莎贝尔 | 奇瓦瓦州 | 1 | 6 |
| 124 | 劳尔 | 瓦哈卡州 | 1 | null |
| 125 | 梅赛德斯 | 瓜达拉哈拉 | 2 | null |
| 126 | 百合 | 坎昆 | 2 | null |
AUTO CDC FROM SNAPSHOT 示例
以下部分提供了用于 AUTO CDC FROM SNAPSHOT 将快照处理到 SCD 类型 1 或类型 2 目标表的示例。 有关何时使用此 API 的背景信息,请参阅 更改数据捕获和快照。
示例:使用管道引入时间处理快照
当快照定期按顺序到达时,可以使用此方法,并依赖管道运行时间戳进行版本控制。 每次管道更新都会获取一个新的快照。
可以从多个源类型(包括 Delta 表、云存储文件和 JDBC 连接)读取快照。
步骤 1:创建示例数据
创建包含快照数据的表。 从管道的 explorations 文件夹中的笔记本或 Databricks SQL 运行以下代码:
CREATE SCHEMA IF NOT EXISTS main.cdc_tutorial;
CREATE TABLE main.cdc_tutorial.snapshot (
userId INT,
city STRING
);
INSERT INTO main.cdc_tutorial.snapshot VALUES
(1, 'Oaxaca'),
(2, 'Monterrey'),
(3, 'Tijuana');
步骤 2:运行 AUTO CDC FROM SNAPSHOT
你需要一个流水线来运行这一步的代码。 请参阅 使用 Lakeflow 管道编辑器开发和调试 ETL 管道。
为快照视图选择源类型(示例创建代码生成 Delta 表):
选项 A:从 Delta 表读取
from pyspark import pipelines as dp
@dp.view(name="source")
def source():
return spark.read.table("main.cdc_tutorial.snapshot")
选项 B:从云存储读取
from pyspark import pipelines as dp
@dp.view(name="source")
def source():
return spark.read.format("csv").option("header", True).load("<snapshot-path>")
选项 C:从 JDBC 读取(仅限经典计算)
from pyspark import pipelines as dp
@dp.view(name="source")
def source():
return (spark.read
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", "<table-name>")
.option("user", "<username>")
.option("password", "<password>")
.load()
)
全部选项,写入到目标
然后添加目标表格和数据流:
dp.create_streaming_table("target")
dp.create_auto_cdc_from_snapshot_flow(
target = "target",
source = "source",
keys = ["userId"],
stored_as_scd_type = 2
)
在第一次管道运行后,所有记录都作为活动行插入:
| userId | city | __START_AT | __END_AT |
|---|---|---|---|
| 1 | 瓦哈卡州 | 0 | null |
| 2 | 蒙特雷 | 0 | null |
| 3 | 提 华纳 | 0 | null |
注释
若要改用 SCD 类型 1 并仅保留当前状态,请设置 stored_as_scd_type=1。 在这种情况下,目标表不包含 __START_AT 和 __END_AT 列。
步骤 3:模拟新快照并重新运行
更新源表以模拟新的快照到达(从管道的 explorations 文件夹中的笔记本或 SQL 文件运行此代码):
TRUNCATE TABLE main.cdc_tutorial.snapshot;
INSERT INTO main.cdc_tutorial.snapshot VALUES
(2, 'Carmel'),
(3, 'Los Angeles'),
(4, 'Death Valley'),
(6, 'Kings Canyon');
再次运行管道。
AUTO CDC FROM SNAPSHOT 将新快照与上一个快照进行比较,并检测到用户 1 已删除、用户 2 和 3 已更新,并插入了用户 4 和 6。 这会生成一个更改提要,并使用 AUTO CDC 来创建输出表。
使用 SCD 类型 2 进行第二次运行后,目标表包含以下记录:
| userId | city | __START_AT | __END_AT |
|---|---|---|---|
| 1 | 瓦哈卡州 | 0 | 1 |
| 2 | 蒙特雷 | 0 | 1 |
| 2 | 卡梅尔 | 1 | null |
| 3 | 提 华纳 | 0 | 1 |
| 3 | 洛杉矶 | 1 | null |
| 4 | 死亡谷 | 1 | null |
| 6 | 国王峡谷 | 1 | null |
用户 1 已结束(已删除)。 用户 2 和 3 各有两个版本显示其城市更改。 新插入了用户 4 和 6。
使用 SCD 类型 1 进行第二次运行后,目标表仅显示当前状态:
| userId | city |
|---|---|
| 2 | 卡梅尔 |
| 3 | 洛杉矶 |
| 4 | 死亡谷 |
| 6 | 国王峡谷 |
示例:使用版本函数处理快照
如果需要显式控制快照排序,请使用此方法。 例如,当多个快照同时到达或快照无序到达时,请使用此方法。 编写一个函数,用于指定要处理的下一个快照及其版本号。 API 按版本号升序处理快照:
- 如果多个快照位于存储中,则它们都按顺序进行处理。
- 如果快照未按顺序到达(例如,
snapshot_3在snapshot_4之后到达),则会跳过它。 - 如果没有新的快照,该函数将
None返回,并且不会执行任何处理。
步骤 1:准备快照文件
创建包含快照数据的 CSV 文件,并将其添加到卷或云存储位置。 按时间顺序命名文件(例如,snapshot_1.csvsnapshot_2.csv)。
每个文件应包含用于 userId 和 city 的列。 例如:
snapshot_1.csv:
| userId | city |
|---|---|
| 1 | 瓦哈卡州 |
| 2 | 蒙特雷 |
| 3 | 提 华纳 |
snapshot_2.csv:
| userId | city |
|---|---|
| 2 | 卡梅尔 |
| 3 | 洛杉矶 |
| 4 | 死亡谷 |
步骤 2:使用版本函数运行 AUTO CDC FROM SNAPSHOT
在流水线中,在文件夹中创建一个新的 Python 代码文件transformations,粘贴以下代码。 然后运行管道。 请参阅 使用 Lakeflow 管道编辑器开发和调试 ETL 管道。
from pyspark import pipelines as dp
from typing import Optional, Tuple
from pyspark.sql import DataFrame
def next_snapshot_and_version(latest_snapshot_version: Optional[int]) -> Optional[Tuple[DataFrame, int]]:
snapshot_dir = "/Volumes/main/cdc_tutorial/snapshots/" # or the location you created the sample data
files = dbutils.fs.ls(snapshot_dir)
snapshot_files = [f.name for f in files if f.name.startswith("snapshot_") and f.name.endswith(".csv")]
snapshot_versions = []
for filename in snapshot_files:
try:
version = int(filename.replace("snapshot_", "").replace(".csv", ""))
snapshot_versions.append(version)
except ValueError:
continue
snapshot_versions.sort()
if latest_snapshot_version is None:
if snapshot_versions:
next_version = snapshot_versions[0]
else:
return None
else:
next_versions = [v for v in snapshot_versions if v > latest_snapshot_version]
if next_versions:
next_version = next_versions[0]
else:
return None
snapshot_path = f"{snapshot_dir}snapshot_{next_version}.csv"
df = spark.read.format("csv").option("header", True).load(snapshot_path)
return (df, next_version)
dp.create_streaming_table("main.cdc_tutorial.target_versioned")
dp.create_auto_cdc_from_snapshot_flow(
target = "main.cdc_tutorial.target_versioned",
source = next_snapshot_and_version,
keys = ["userId"],
stored_as_scd_type = 2
)
注释
若要改用 SCD 类型 1,请设置 stored_as_scd_type=1。
处理 snapshot_1.csv后,目标表包含以下记录:
| userId | city | __START_AT | __END_AT |
|---|---|---|---|
| 1 | 瓦哈卡州 | 1 | null |
| 2 | 蒙特雷 | 1 | null |
| 3 | 提 华纳 | 1 | null |
处理 snapshot_2.csv后,目标表包含以下记录:
| userId | city | __START_AT | __END_AT |
|---|---|---|---|
| 1 | 瓦哈卡州 | 1 | 2 |
| 2 | 蒙特雷 | 1 | 2 |
| 2 | 卡梅尔 | 2 | null |
| 3 | 提 华纳 | 1 | 2 |
| 3 | 洛杉矶 | 2 | null |
| 4 | 死亡谷 | 2 | null |
注释
请记住,对于 SCD 类型 1,表看起来与最新的快照完全相同。 区别在于,下游查询可以使用更改源来仅处理更改的记录。
步骤 3:添加新快照
使用修改后的数据(例如,已更改的城市值、新行或删除行)将新的 CSV 文件添加到存储位置。 然后再次运行管道以处理新快照。
局限性
- 排序列必须是可排序的数据类型。
NULL不支持序列化值。 -
AUTO CDC FROM SNAPSHOT仅在 Python 管道接口中受支持;不支持 SQL 接口。 - 若要从 AUTO CDC 过程的目标端流式读取数据,请从其更改馈送中读取。 有关详细信息,请参阅 从 AUTO CDC 目标表读取变更数据流。
其他资源
- 更改数据捕获和快照:了解 CDC 概念、快照和 SCD 类型。
-
使用
AUTO CDC复制外部 RDBMS 表:了解如何使用once流执行初始注入,然后继续处理更改。 - 高级 AUTO CDC 主题:了解 AUTO CDC 目标的更改操作、读取更改数据馈送和处理指标。
- 倒带并重放管道:了解如何将 AUTO CDC 目标还原到较早的时间点并重新处理变更。
- 教程:使用变更数据捕获生成 ETL 管道