AUTO CDC API:使用管道简化变更数据捕获

Lakeflow 管道借助 AUTO CDCAUTO 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_flowcreate_auto_cdc_from_snapshot_flow

注释

本页介绍如何根据源数据的变化来更新管道中的表。 若要了解如何记录和查询 Delta 表的行级更改信息,请参阅在Azure Databricks上使用更改数据馈送

要求

若要使用 CDC API,必须将您的管道配置为使用 无服务器 Lakeflow 管道 或 Lakeflow 管道 ProAdvanced版本

AUTO CDC 的工作原理

若要执行 AUTO CDC CDC 处理,请首先创建流式处理表,然后分别在 SQL 中使用 AUTO CDC ... INTO 语句或在 Python 中使用 create_auto_cdc_flow() 函数来指定变更跟踪源、关键字和排序。 有关序列化和 SCD 逻辑的工作原理的说明,请参阅 更改数据捕获和快照。 请参阅 AUTO CDC 示例

若要从具有更改源的源进行初始数据注入,请使用 AUTO CDConce 流,然后继续处理更改源。 请参阅 使用 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

注释

以下所有示例都包含指定DELETETRUNCATE操作的选项,但每个操作都是可选的。

创建示例数据

运行以下语句以创建示例数据集。 此代码不应作为管道定义的一部分运行。 从管道的浏览文件夹运行它,而不是转换文件夹。

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 覆盖以前的值。 之前在 UPDATEsequenceNum=5 被删除,因为在 sequenceNum=6 的更新已经到来。

在运行未注释的 TRUNCATE 记录的示例后,表在 sequenceNum=3 处被清除。 这意味着记录 124126 不在表中,最终目标表仅包含以下记录:

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_3snapshot_4 之后到达),则会跳过它。
  • 如果没有新的快照,该函数将 None 返回,并且不会执行任何处理。

步骤 1:准备快照文件

创建包含快照数据的 CSV 文件,并将其添加到卷或云存储位置。 按时间顺序命名文件(例如,snapshot_1.csvsnapshot_2.csv)。

每个文件应包含用于 userIdcity 的列。 例如:

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 目标表读取变更数据流

其他资源