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 还支持部分更新,其中更改记录仅更新列的子集。 请参阅 “应用部分更新”。

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

使用多个列进行排序

若要按多个列(例如时间戳和 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 仅保留每个记录的最新版本。 以下示例从上面创建的更改数据馈送中读取数据,并将更改应用于流处理表目标。 你需要一个流水线来运行这些代码。 请参阅 使用 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 覆盖以前的值。 之前在 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

你需要一个流水线来运行这一步的代码。 请参阅 使用 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_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

在流水线中,在文件夹中创建一个新的 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 目标表读取变更数据流

其他资源