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 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。
双时态 AUTO CDC 的工作原理
Important
Bitemporal AUTO CDC 处于 Beta 阶段。
SCD 类型 1 和类型 2 都是单时态的:它们仅沿一个时间维度跟踪变更。 双时态在 SCD 类型 2 历史记录的基础上进行扩展,用于跟踪两个时间维度上的变化,并区分两种视角:
- 业务时间:事件发生时。
- 系统时间:系统记录或引入事件的时间。
与 SCD 类型 2 一样,bitemporal 保留完整的记录历史记录。 它增加了第二条时间线,因此你可以还原过去任意时间点上数据当时显示的内容以及系统当时的判断。
例如,对冲基金从源系统引入股票数据。 Acme Corp 的股价于1月1日发生变化,但该基金直到1月5日才纳入这一更新。 双时态 AUTO CDC 使基金能够回答两个截然不同的问题:Acme Corp 在 1 月 1 日的实际股价是多少(业务时间),以及该基金在 1 月 3 日作出交易决策时,系统认为的价格是多少(系统时间)。 区分这些时间线的能力对于审核、监管报告和财务决策非常有用。
若要启用双时间处理,请设置 STORED AS BITEMPORAL(SQL)或 stored_as_scd_type="bitemporal"(Python),将 SEQUENCE BY 用于业务时间列,并将 SYSTEM SEQUENCE BY 用于系统时间列。 目标表添加了 __SYSTEM_START_AT 列和 __SYSTEM_END_AT 列,以及 SCD 类型 2 的 __START_AT 列和 __END_AT 列。 有关语法详细信息,请参阅 AUTO CDC INTO(管道)或 create_auto_cdc_flow。
有关插入、更新、乱序更新和删除如何影响双时态表的分步讲解,请参阅双时态 AUTO CDC 示例。
使用多个列进行排序
若要按多个列(例如时间戳和 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 仅保留每个记录的最新版本。 以下示例从上面创建的更改数据馈送中读取数据,并将更改应用于流处理表目标。 什么是管道? 以运行此代码。
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
使用什么是流水线?来运行此步骤中的代码。
为快照视图选择源类型(示例创建代码生成 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
创建新的笔记本并粘贴以下管道代码。 那么什么是管道?
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 文件添加到存储位置。 然后再次运行管道以处理新快照。
双时态 AUTO CDC 示例
Important
Bitemporal AUTO CDC 处于 Beta 阶段。
以下示例从一小组合成的 CDC 事件创建一个双时态目标表。 该 bt 列具有业务时间, st 该列具有系统时间。
Python
from pyspark import pipelines as dp
# Source: synthetic CDC events
dp.create_streaming_table(name="cdc_source")
@dp.append_flow(target="cdc_source", once=True)
def load_cdc_source():
return spark.createDataFrame(
[
(1, "x10", "y10", 10, 100),
(1, "x20", "y20", 20, 200)
],
schema="id INT, x STRING, y STRING, bt INT, st INT",
)
# Target: bitemporal table
dp.create_streaming_table(name="target_bitemporal")
dp.create_auto_cdc_flow(
target = "target_bitemporal",
source = "cdc_source",
keys = ["id"],
sequence_by = "bt",
system_sequence_by = "st",
stored_as_scd_type = "bitemporal"
)
SQL
-- Source: synthetic CDC events
CREATE OR REFRESH STREAMING TABLE cdc_source_sql;
CREATE FLOW cdc_source_sql AS INSERT INTO ONCE
cdc_source_sql BY NAME
SELECT * FROM VALUES
(1, 'x10', 'y10', 10, 100),
(1, 'x20', 'y20', 20, 200)
AS t(id, x, y, bt, st);
-- Target: bitemporal table
CREATE OR REFRESH STREAMING TABLE target_bitemporal_sql;
CREATE FLOW target_bitemporal_sql AS AUTO CDC INTO
target_bitemporal_sql
FROM
stream(cdc_source_sql)
KEYS
(id)
SEQUENCE BY
bt
SYSTEM SEQUENCE BY
st
STORED AS
BITEMPORAL;
以下步骤将逐步演示双时态表如何针对单个公司记录插入、更新、乱序更新和删除操作。 排序列生成 __START_AT 和 __END_AT (业务时间)列,系统排序列生成 __SYSTEM_START_AT 和 __SYSTEM_END_AT (系统时间)列:
| 列 | Description |
|---|---|
__START_AT |
此行生效的业务时间。 |
__END_AT |
此行有效性结束的业务时间。
null 如果永久有效。 |
__SYSTEM_START_AT |
此行数据及其业务时间区间被认定为真实有效时的系统时间。 |
__SYSTEM_END_AT |
此行数据及其业务时间区间被认定为失效时的系统时间。
null 如果已知将始终为 true。 |
系统处理在这两个时间线中按任意顺序到达的事件。 当某个事件到达时,如果其业务时间或系统时间早于已处理过的事件,系统会更正受影响的历史记录,而不是仅追加到末尾。
步骤 1:插入
公司 A 于 2025 年 7 月 18 日 10:01:00(业务时间)被添加,但直到 10:05:00(系统时间)才被导入。
输入:
| CompanyId | 数据点 | 先后顺序 | 系统排序 | 运算 |
|---|---|---|---|---|
| A | XFv1 | 7/18/2025 10:01:00 | 7/18/2025 10:05:00 | INSERT |
输出:
| CompanyId | 数据点 | __START_AT | __END_AT | __SYSTEM_START_AT | __SYSTEM_END_AT |
|---|---|---|---|---|---|
| A | XFv1 | 7/18/2025 10:01:00 | Null | 7/18/2025 10:05:00 | Null |
XFv1 自 10:01:00 起有效,结束时间未知。 系统在系统时间 10:05:00 获知这一事实,且结束时间未知。
步骤 2:更新
公司 A 于 2025 年 7 月 18 日 12:15:43(业务时间)更新,系统于 12:20:00(系统时间)处理该事件。 系统会同时保留两部分内容:一是得知该更新之前其原先认定的内容,二是导入该更新后更正后的业务历史。
输入:
| CompanyId | 数据点 | 先后顺序 | 系统排序 | 运算 |
|---|---|---|---|---|
| A | XFv2 | 7/18/2025 12:15:43 | 7/18/2025 12:20:00 | UPDATE |
输出:
| CompanyId | 数据点 | __START_AT | __END_AT | __SYSTEM_START_AT | __SYSTEM_END_AT |
|---|---|---|---|---|---|
| A | XFv1 | 7/18/2025 10:01:00 | Null | 7/18/2025 10:05:00 | 7/18/2025 12:20:00 |
| A | XFv1 | 7/18/2025 10:01:00 | 7/18/2025 12:15:43 | 7/18/2025 12:20:00 | Null |
| A | XFv2 | 7/18/2025 12:15:43 | Null | 7/18/2025 12:20:00 | Null |
XFv1 被认为自 10:01:00 起有效,且没有已知的结束时间,系统从 10:05:00 一直这样认为,直到 12:20:00。 现已知 XFv1 仅有效至 12:15:43;一条自系统时间 12:20:00 起生效的更正历史记录目前没有已知结束时间。 XFv2 自 12:15:43 起有效,结束时间未知,并于系统时间 12:20:00 被获知。
步骤 3:乱序更新
收到一条乱序更新,表明 A 公司实际上已于 2025 年 7 月 18 日 12:05:00(业务时间)更新,但直到 12:25:00(系统时间)才被引入。 当某个更新在系统时间上较晚到达,但其业务时间却早于先前记录的业务时间时,系统会修正历史业务时间,并同时保留系统在该乱序更新到达之前所认为的内容以及更正后的历史记录。
输入:
| CompanyId | 数据点 | 先后顺序 | 系统排序 | 运算 |
|---|---|---|---|---|
| A | XFv3 | 7/18/2025 12:05:00 | 7/18/2025 12:25:00 | UPDATE |
输出:
| CompanyId | 数据点 | __START_AT | __END_AT | __SYSTEM_START_AT | __SYSTEM_END_AT |
|---|---|---|---|---|---|
| A | XFv1 | 7/18/2025 10:01:00 | Null | 7/18/2025 10:05:00 | 7/18/2025 12:20:00 |
| A | XFv1 | 7/18/2025 10:01:00 | 7/18/2025 12:15:43 | 7/18/2025 12:20:00 | 7/18/2025 12:25:00 |
| A | XFv1 | 7/18/2025 10:01:00 | 7/18/2025 12:05:00 | 7/18/2025 12:25:00 | Null |
| A | XFv3 | 7/18/2025 12:05:00 | 7/18/2025 12:15:43 | 7/18/2025 12:25:00 | Null |
| A | XFv2 | 7/18/2025 12:15:43 | Null | 7/18/2025 12:20:00 | Null |
XFv1 被认为在 10:01:00 到 12:15:43 期间有效,并且这一判断现在在系统时间上有效至 12:25:00。 此次更新将 XFv1 的业务有效期更正为于 12:05:00 结束,更正后的历史记录自系统时间 12:25:00 起生效。 现已知 XFv3 在 12:05:00 至 12:15:43 期间有效;这一认知在系统时间自 12:25:00 起成立,且结束时间未知。
步骤 4:删除
公司 A 于 2025 年 7 月 18 日 12:30:00 被删除,系统于 12:30:00 处理该事件。 由于删除操作表示实体业务存在的末尾,因此系统不会创建替换行。 XFv2 分两行显示,同时保留了完整的审计跟踪记录,既包括该公司何时不复存在,也包括系统何时获知该删除操作。
输入:
| CompanyId | 数据点 | 先后顺序 | 系统排序 | 运算 |
|---|---|---|---|---|
| A | XFv2 | 7/18/2025 12:30:00 | 7/18/2025 12:30:00 | DELETE |
输出:
| CompanyId | 数据点 | __START_AT | __END_AT | __SYSTEM_START_AT | __SYSTEM_END_AT |
|---|---|---|---|---|---|
| A | XFv1 | 7/18/2025 10:01:00 | Null | 7/18/2025 10:05:00 | 7/18/2025 12:20:00 |
| A | XFv1 | 7/18/2025 10:01:00 | 7/18/2025 12:15:43 | 7/18/2025 12:20:00 | 7/18/2025 12:25:00 |
| A | XFv1 | 7/18/2025 10:01:00 | 7/18/2025 12:05:00 | 7/18/2025 12:25:00 | Null |
| A | XFv3 | 7/18/2025 12:05:00 | 7/18/2025 12:15:43 | 7/18/2025 12:25:00 | Null |
| A | XFv2 | 7/18/2025 12:15:43 | Null | 7/18/2025 12:20:00 | 7/18/2025 12:30:00 |
| A | XFv2 | 7/18/2025 12:15:43 | 7/18/2025 12:30:00 | 7/18/2025 12:30:00 | Null |
XFv2 自 12:15:43 起有效,且无已知的结束时间,系统在 12:20:00 至 12:30:00 期间一直认为如此。 该删除操作被摄入后,已知 XFv2 仅在 12:30:00 之前有效;更正后的历史记录自系统时间 12:30:00 起生效。
局限性
- 排序列必须是可排序的数据类型。
NULL不支持序列化值。 -
AUTO CDC FROM SNAPSHOT仅在 Python 管道接口中受支持;不支持 SQL 接口。 - 若要从 AUTO CDC 过程的目标端流式读取数据,请从其更改馈送中读取。 有关详细信息,请参阅 从 AUTO CDC 目标表读取变更数据流。
其他资源
- 更改数据捕获和快照:了解 CDC 概念、快照和 SCD 类型。
-
使用
AUTO CDC复制外部 RDBMS 表:了解如何使用once流执行初始注入,然后继续处理更改。 - 高级 AUTO CDC 主题:了解 AUTO CDC 目标的更改操作、读取更改数据馈送和处理指标。
- 教程:使用变更数据捕获生成 ETL 管道