Databricks

Taking AUTO CDC to the next level: Solving the hardest real-world use cases

8.5内容质量

TL;DR · AI 摘要

Databricks在Apache Spark 4.2中扩展AUTO CDC功能,支持Bitemporal和Partial Updates,解决数据工程中的复杂问题。

核心要点

  • Bitemporal CDC通过双时间轴(业务时间+系统时间)满足SEC合规要求
  • Partial Updates支持安全处理缺失字段,避免数据损坏
  • AUTO CDC能力已开源至Apache Spark 4.2

结构提纲

按章节快速跳转。

  1. 介绍AUTO CDC在Spark中的演进及新功能价值

  2. 通过双时间轴实现SEC合规要求的审计追踪

  3. 安全处理字段缺失的变更数据捕获方案

  4. 开源版本新增乱序处理和审计能力

思维导图

用一张图看清主题之间的关系。

查看大纲文本(无障碍 / 无 JS 友好)
  • AUTO CDC演进
    • Bitemporal CDC
      • 双时间轴管理
      • SEC合规支持
    • Partial Updates
      • 字段级变更处理
      • 数据完整性保障
    • Spark 4.2扩展
      • 开源生态支持
      • 乱序处理能力

金句 / Highlights

值得收藏与分享的关键句。

#Spark#CDC#数据工程#Apache Spark 4.2
打开原文

将 AUTO CDC 推向更高层次:解决最复杂的实际用例 | Databricks 博客

跳至主要内容

产品

2026年8月11日

将 AUTO CDC 推向更高层次:解决最复杂的实际用例

从双时间合规到部分记录更新:无需自定义代码即可实现强大且可审计的变更数据捕获

作者:Josh Seidel、Shanelle Roman 和 Sudhanva Huruli

摘要

  • AUTO CDC 用声明式管道取代手动编写的 MERGE 逻辑,实现变更数据捕获
  • Spark 声明式管道现在支持双时间 AUTO CDC,可独立跟踪业务时间和系统时间,同时支持部分更新以安全处理缺失字段
  • AUTO CDC 功能已扩展到开源 Apache Spark 4.2,为更广泛的生态系统带来标准化的乱序变更数据捕获

变更数据捕获是数据工程师在 Spark 上构建的最常见功能之一,但手动实现这一功能却是最繁琐的。在我们之前的博文《停止手动编写变更数据捕获管道》中,我们介绍了 Apache™ Spark 声明式管道(SDP)中的 AUTO CDC 如何通过几个简单的声明取代数百行脆弱的 MERGE 逻辑,实现 SCD 类型 1、SCD 类型 2 和快照 CDC 的自动化。

随着管道需求的演变,工程师会遇到标准 CDC 模式难以解决的情况:

  • 处理无序的双时间时间线
  • 在不破坏现有数据的情况下处理部分记录更新
  • 保留超出存储保留窗口的可审计性

今天,我们通过提升 AUTO CDC 的能力来解决这些实际挑战,并将这些功能扩展到开源 Apache Spark 4.2。

使用双时间 AUTO CDC 进行双轴历史跟踪

标准的 SCD 类型 2 表可以告诉你现实世界中事实变更的时间,但无法告诉你系统在任意时间点的信念状态。

根据 SEC 第 17a-4 条规则和 FINRA 记录保存规则,企业必须能够重建特定时间点的记录;仅 SEC 的记录保存审查就已使 100 多家企业自 2021 年以来累计被罚款超过 20 亿美元。困难之处通常不在于存储当前值,而在于数月后回答:在报告日期参考数据说了什么?我们的系统当时相信什么?

标准的 SCD 类型 2 仅跟踪一个时间线:事实变更的时间。双时间 AUTO CDC 独立跟踪两个时间线:

  • 业务时间(又称事件时间或有效时间):事实在现实世界中实际为真的时间。例如,股票代码在周一变为可报告项;国家代码在季度末被弃用。
  • 系统时间(又称事务时间或处理时间):系统记录学习到数据的时间。周一的变更可能要到周三才会进入管道。

每个目标表会获得四列系统管理列:__START_AT 和 __END_AT 用于业务时间,__SYSTEM_START_AT 和 __SYSTEM_END_AT 用于系统时间。单个逻辑事实可以有多个物理行,每个业务版本/系统版本组合对应一行,这使得沿任一轴进行时间点重建成为可能。关键行为保证:事件可以按任意顺序出现在任一时间线中。

当更正记录的业务时间或系统时间早于已处理的数据时,引擎会重写受影响的历史记录,而不是简单地追加到末尾。无需手动编写逻辑,只需声明两个排序列,引擎即可维护两个时间间隔。这种机制同样适用于维度表(如符号主表)和事实表(如交易历史或传感器读数),这些表需要严格的可审计性。以下是其在FINRA CAT参考数据中的应用示例:

$

/$

请注意,确切的SQL子句是STORED AS BITEMPORAL,而非STORED AS SCD TYPE BITEMPORAL,且需要同时指定SEQUENCE BY和SYSTEM SEQUENCE BY。假设Acme的可报告标志在1月1日(业务时间)变更,但数据源直到1月5日(系统时间)才接收到该变更。随后在1月8日收到一份回溯更正,称实际变更发生在1月1日但数值不同。双时间AUTO CDC可以回答这两个问题:

1月3日,第一个查询返回空结果,这是系统在当时显示的正确可审计答案。第二个查询在今天执行时会反映更正后的真相。两个时钟对应两个答案,两者都正确。排序列必须为可排序类型,且不允许存在NULL排序值。该功能可在无服务器SDP或Pro/Advanced产品版本上运行,目前处于Beta阶段,因此请将管道固定到渠道:PREVIEW。

超越时间旅行:可重复的机器学习在VACUUM后依然有效

当模型在参考数据或特征数据上训练时,可重复性意味着在数月后的审查或审计中,能够重建模型当时使用的精确数据集。直觉上人们会转向Delta Lake时间旅行功能,但这是表文件历史的属性,而非永久记录。VACUUM会永久删除不再被近期版本引用的数据文件;一旦超过默认的7天保留窗口,训练时记录的TIMESTAMP AS OF可能突然失效。双时间表将历史作为数据存储,而非文件版本。VACUUM和OPTIMIZE会压缩文件但不会影响逻辑历史,因此所有过去的业务版本或系统版本仍然是可查询的行。获取可重复性的两种方式是:将两个时间点(业务时间和系统时间)作为MLflow参数记录,并将训练查询固定到该信念状态:

或者,如果表暴露了当前视图,可在训练时记录一个系统时间点,并在之后通过该时间点的系统时间查询重建数据:

无论哪种方式,可重复性契约只需在MLflow运行中记录两个时间戳。由于双时间历史以行形式存储,即使VACUUM清理了底层文件,该契约依然有效。

AutoCDC部分更新功能现已正式发布

并非所有变更数据捕获(CDC)源在更新时都会发送完整行。相反,许多源仅发送已更改的字段,将其他列表示为NULL。如果没有特殊处理,这些NULL值可能会无意中覆盖目标表中的现有数据。在此之前,客户必须构建自定义逻辑来处理这种行为。现在,AutoCDC部分更新功能可自动处理这种情况。

部分更新通过允许更新事件仅修改部分列来扩展AutoCDC功能。对于选定的列,传入更新中的NULL值会被解释为“不更新”,而非覆盖现有值。

部分更新功能对于处理 CDC 数据源中因发出 NULL 而省略未更改值的情况特别有用。如果没有启用部分更新,这些 NULL 值会覆盖目标表中已有的数据。

例如,假设目标表包含:(1, 'A', 20)

一个传入的更新事件包含:(1, NULL, 30)

默认情况下,AutoCDC 会将该行更新为:(1, NULL, 30)。

启用部分更新后,名称字段中的 NULL 会被视为"保留现有值不变",最终结果为:(1, 'A', 30)。

启用部分更新只需在 AutoCDC 定义中添加一个参数。您可以选择以下三种方式指定哪些列应作为部分更新处理:

  • 忽略 NULL 值的列列表:IGNORE NULL UPDATES ON columnList
  • 不忽略 NULL 值的列列表:IGNORE NULL UPDATES ON * EXCEPT (columnList)
  • 每行可具有不同源列名的更新列:COLUMNS TO UPDATE

如需完整语法、示例和使用指南,请参阅《应用部分更新》文档。

我们继续致力于开源

Spark Declarative Pipelines 是开源的,因此其最广泛使用的流类型也应如此。我们首先将 AUTO CDC Type 1 的 Python API 贡献给 Apache Spark 4.2。

我们采用 Spark 的演进方式来贡献代码:通过一系列经过审查的提案和拉取请求,而非一次性代码提交(参见 SPIP 和 SPARK-56249)。乱序数据的正确性已内置:一个小型辅助表会跟踪早期到达事件(如删除墓碑)的状态,重试的微批次会收敛而非破坏目标表,因为它基于 Spark 的流处理和表抽象而非存储格式,因此可在 Delta Lake 和 Apache Iceberg 上运行。

开源路线图中的下一步:

  • 下个版本功能:我们已将 SQL 接口(CREATE FLOW ... AS AUTO CDC INTO)合并到主分支,该功能将在下一个 Apache Spark 版本中发布。
  • 高级管道语义:正在开发 SCD Type 2 完整历史管理、原生变更日志输入以及部分更新支持,以防止 NULL 值覆盖目标数据。
  • 可靠性与测试:我们正在添加"应用即截断"功能,同时扩展围绕乱序数据和幂等重试的自动化测试套件。

入门指南

无论您是要实现双时间合规性、设置部分更新,还是探索 Apache Spark 中的开源 AutoCDC,都可以查看以下资源快速上手:

  • 双时间 AUTO CDC:学习如何在 SDP 中配置双轴历史跟踪
  • 部分更新指南:查看完整语法和示例,了解如何在不使用自定义代码的情况下处理缺失的更新字段
  • 开源 AutoCDC 编程指南:探索 Apache Spark 4.2 中可用的声明式 CDC 功能

订阅最新文章

订阅我们的博客,即可将最新文章直接发送到您的邮箱。

注册订阅

查看所有博客

slice-start id="_gatsby-scripts-1"

slice-end id="_gatsby-scripts-1"