Taking AUTO CDC to the next level: Solving the hardest real-world use cases
TL;DR · AI 摘要
Databricks在Apache Spark 4.2中扩展AUTO CDC功能,支持Bitemporal和Partial Updates,解决数据工程中的复杂问题。
核心要点
- Bitemporal CDC通过双时间轴(业务时间+系统时间)满足SEC合规要求
- Partial Updates支持安全处理缺失字段,避免数据损坏
- AUTO CDC能力已开源至Apache Spark 4.2
结构提纲
按章节快速跳转。
思维导图
用一张图看清主题之间的关系。
查看大纲文本(无障碍 / 无 JS 友好)
- AUTO CDC演进
- Bitemporal CDC
- 双时间轴管理
- SEC合规支持
- Partial Updates
- 字段级变更处理
- 数据完整性保障
- Spark 4.2扩展
- 开源生态支持
- 乱序处理能力
金句 / Highlights
值得收藏与分享的关键句。
SEC的记录检查已导致2021年后100+公司累计超20亿美元罚款
Bitemporal CDC维护__START_AT/__END_AT等4个系统管理列
Partial Updates避免因缺失字段导致现有数据损坏
将 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"