将AUTO CDC提升至新高度:解决最棘手的现实世界用例

将AUTO CDC提升至新高度:解决最棘手的现实世界用例

💡 原文英文,约1200词,阅读约需5分钟。
📝

内容提要

本文介绍Apache Spark Declarative Pipelines中AUTO CDC的增强功能。新增双时态CDC,通过系统时间和业务时间双时间轴,支持任意顺序事件到达及历史修正,满足SEC审计要求。部分更新功能自动处理NULL值,避免覆盖现有数据。此外,AUTO CDC Type 1已贡献至开源Apache Spark 4.2,支持Delta Lake和Iceberg。

🔎

延伸解读

双时态CDC的审计价值

双时态CDC通过系统时间和业务时间两条独立时间轴,解决了标准SCD Type 2无法回答“系统在特定时间点相信什么”的问题。文章指出,SEC规则要求企业重建历史记录,且自2021年以来相关罚款已超20亿美元。该功能允许事件乱序到达,并自动重写历史,无需手写逻辑,适用于维度表和事实表,满足严格审计要求。

历史存储为数据而非文件版本

双时态表将历史存储为数据行,而非依赖Delta Lake的时间旅行。VACUUM清理文件后,时间旅行查询可能失效,但双时态表仍可通过查询行重建历史。这为MLflow模型训练的可复现性提供了更可靠的保障,只需记录业务和系统时间戳,即可在审计或复核时重建训练数据集。

部分更新避免NULL覆盖

许多CDC源只发送变更字段,未变更字段以NULL表示,若直接处理会覆盖现有数据。AUTO CDC部分更新功能自动将NULL解释为“不更新”,支持三种指定列的方式:IGNORE NULL UPDATES ON、IGNORE NULL UPDATES ON * EXCEPT、COLUMNS TO UPDATE。这消除了手写自定义逻辑的需求,简化了管道维护。

开源贡献与跨格式支持

AUTO CDC Type 1已贡献给Apache Spark 4.2,通过辅助表处理乱序事件,确保正确性。它基于Spark的流和表抽象,不依赖特定存储格式,因此同时支持Delta Lake和Iceberg。这为开源社区提供了更广泛的CDC能力,降低了构建复杂管道的门槛。

Q&A

什么是双时态CDC?它解决了什么问题?

双时态CDC是AUTO CDC的一项新功能,通过系统时间和业务时间两个独立时间轴跟踪数据变化。它解决了标准SCD Type 2只能回答事实何时变化,而无法回答系统在特定时间点所相信的状态的问题,满足SEC和FINRA的审计要求。

双时态CDC如何实现点时间重建?

双时态CDC通过四个系统管理列(__START_AT、__END_AT、__SYSTEM_START_AT、__SYSTEM_END_AT)实现,每个逻辑事实可以有多个物理行,对应不同的业务时间和系统时间组合。这样,可以沿任一轴进行时间点重建,例如查询系统在特定日期的状态或修正后的真实状态。

双时态CDC在哪些场景下适用?

双时态CDC适用于需要严格审计性的维度表和事实表,如FINRA CAT参考数据、交易历史或传感器读数。它能够回答系统在特定时间点的状态以及修正后的真实状态。

双时态CDC如何保证事件乱序到达时的正确性?

双时态CDC允许事件在任一时间轴上任意顺序到达。当出现更早业务时间或系统时间的修正时,引擎会重写受影响的历史记录,而不是简单追加。用户只需声明两个排序列,引擎自动维护两个时间间隔。

双时态CDC的Beta版本有哪些使用限制?

双时态CDC目前处于Beta阶段,需要在SDP上使用PREVIEW频道,并要求排序列必须是可排序类型且无NULL值。

双时态表如何支持模型训练的可复现性?

双时态表将历史存储为数据而非文件版本,即使VACUUM清理文件,历史仍可通过查询行重建。可以通过记录MLflow参数(业务时间和系统时间)来固定训练时的数据状态,实现可复现性。

AUTO CDC部分更新功能是什么?它如何处理NULL值?

AUTO CDC部分更新功能允许更新事件只修改部分列,对于选定的列,NULL值被解释为“不更新”,而不是覆盖现有值。这避免了CDC源中未变化的列以NULL形式发送时覆盖目标表中的现有数据。

如何指定部分更新的列?

有三种方式:IGNORE NULL UPDATES ON columnList(忽略指定列的NULL更新)、IGNORE NULL UPDATES ON * EXCEPT (columnList)(忽略除指定列外的所有列的NULL更新)、COLUMNS TO UPDATE(指定要更新的列)。

AUTO CDC Type 1在开源Apache Spark 4.2中如何支持Delta Lake和Iceberg?

AUTO CDC Type 1已贡献给Apache Spark 4.2,它基于Spark的流和表抽象,而非特定存储格式,因此同时支持Delta Lake和Iceberg。它通过辅助表处理乱序事件,确保正确性。

🏷️

标签

➡️

继续阅读