Databricks Delta Lake中流检查点与变更流有何差异?
Databricks Delta Lake:流检查点(Streaming Checkpoint)与变更流(Change Stream)的核心差异
1. 流检查点(Streaming Checkpoint)
流检查点是流作业的状态管理机制,核心作用是保障流处理的可靠性:
- 记录流作业的运行进度(比如Delta表的版本号、Kafka的偏移量)、窗口聚合状态、数据连接状态等关键信息。
- 作业运行时会定期将这些状态写入指定的存储目录(通常是云存储路径),当作业故障重启时,自动读取检查点恢复到断点,避免重复计算或数据丢失,实现Exactly-Once/At-Least-Once语义。
- 示例配置:
spark.readStream .format("delta") .load("/path/to/source-table") .writeStream .format("delta") .option("checkpointLocation", "/path/to/checkpoint-dir") # 指定检查点目录 .start("/path/to/sink-table")
2. 变更流(Change Stream)
变更流是Delta Lake内置的增量数据捕获(CDC)能力,用来追踪并消费Delta表的所有变更:
- 可以捕获Delta表中发生的插入、更新、删除、合并等操作,输出包含操作类型(如
insert/update/delete)、变更前后数据的事件流。 - 支持从指定版本开始回溯消费历史变更,或实时监听新产生的变更,常用于构建数据同步管道、审计日志、实时数仓更新等场景。
- 示例读取变更流:
spark.readStream .format("delta") .option("readChangeFeed", "true") # 开启变更流读取 .option("startingVersion", "5") # 从版本5开始消费 .load("/path/to/delta-table") .writeStream .format("console") .start()
3. 核心差异对比
- 定位与用途:
- 流检查点:是流作业的可靠性保障工具,解决的是作业故障恢复、语义一致性问题,不直接处理业务数据变更。
- 变更流:是数据变更的消费接口,解决的是如何获取Delta表的增量变更,服务于下游业务逻辑。
- 依赖关系:
- 两者可以配合使用:当用流作业消费变更流时,必须配置检查点,保证作业重启后能从上次消费的位置继续,避免重复处理变更事件。
- 存储内容:
- 流检查点存储的是作业状态数据(偏移量、聚合状态等),与业务数据无关。
- 变更流输出的是业务数据的变更事件,直接关联业务逻辑。
- 运行逻辑:
- 流检查点是流作业后台自动维护的,用户只需要指定存储路径即可,无需主动操作。
- 变更流需要用户主动发起读取请求,以流或批的方式获取变更数据。
内容的提问来源于stack exchange,提问作者ng.newbie
相关产品推荐
相关产品推荐

