回退和重播管道

Important

此功能在 Beta 版中。 工作区管理员可以从 预览 页控制对此功能的访问。 请参阅 Manage Azure Databricks 预览版

错误的转换、源记录批次的错误或意外的模式变更,都可能导致管道表从已知时间点开始出现错误。 倒带会将流水线返回问题发生前的某个点,这样你可以部署修复并只重新处理受影响的数据。

Rewind 会同时恢复表版本、流式数据源偏移量和算子状态,从而确保回放时不会跳过记录或写入重复记录。 三种操作对数据进行重处理,解决不同的问题:

  • 倒带适用于可恢复的流水线,这类流水线自某个已知时间点起写入了错误数据,例如在错误转换、格式错误的输入或错误的代码部署之后。 它会将表数据、源偏移量和操作符状态恢复到问题发生前的某个点,并仅重新处理受影响的数据,保持操作符状态。 该时间点之前的数据未被修改。
  • 完全刷新使用所有可用的源数据重新构建表,并丢弃其当前内容。 用它从头重新计算,或者当代码变更与现有状态不兼容时。
  • 检查点重置 用于恢复检查点无效或已损坏,或者因与检查点不兼容的代码变更而被阻塞的流水线。 它会重置检查点,并在保留表的当前内容的同时继续向前推进。 Rewind无法恢复这些情况,因为它没有放宽结构化流兼容性规则。

要求

Requirement 详情
通道 管道必须位于 预览 通道。 请参阅 “配置管道”。
Configuration 将流水线配置true设置为 pipelines.rewind.betaEnabled ,然后运行一次流水线。 每个流程只有在完成一次启用了时间旅行的更新后,才会变为可回退。
管道模式 触发式和连续管道。 不支持实时模式。
Sources Delta 表、流式表、Kafka 和 Auto Loader。
目标 流式表格和实体视图。
Flows 流式数据流和 AUTO 变更数据捕获(CDC)流,包括 SCD 类型 1 和 SCD 类型 2 目标表。 支持有状态查询,例如聚合、联接和去重。

管道中的每一次流量都必须满足这些要求。 当 pipelines.rewind.betaEnabledtrue 时,包含不符合条件的流的管道更新将会失败。 在启用之前,请确认每个流程都满足上述要求。

注释

pipelines.rewind.betaEnabled 设置为 true 的流水线无法移回 Current 通道,除非 Current 通道更新到支持回退的运行时。

倒带和回放的工作原理

倒带和回放是两个独立的步骤。

倒带会将每个表恢复到其在倒带点时的版本,并重置用于记录每个流已读取到何处的流式检查点。 你的转换不会运行,也不会重新处理源数据。

重放会在流水线下次运行时发生。 它会按照当前的流水线定义从回退点开始重新处理,追赶至当前状态,然后恢复正常的增量处理。 Rewind 不会启动该流程,因此请在准备就绪后自行启动。

该管道会自动生成回退点,频率约为每小时一次,并将其保留 7 天。 你刚创建的管道在产生第一个结果之前,还没有可回退到的内容。

回退某个数据集也会将同一流水线中其下游的所有内容一并回退。 回退支持一个管道。 它不与其他管道或同一表的外部读取器协调,因此应分别管理这些数据。

使用 UI 回退管道

使用 Rewind 的两种主要方式是 UI 和 Genie。 界面会列出可用的回退点,并在你提交前显示每个回退点会影响哪些数据集。

  1. 在管道页面上,点击 按钮(位于 下倒V字形图标。运行管道 旁边),然后点击 回退管道
  2. 选择一个倒带点,或者使用快捷方式,比如“ 倒带到昨天 ”或“ 倒带到最新时间点”。 单击 “下一步”
  3. 选择包含哪些表格。 使用 视图选择管道图中的数据集,或使用 列表 视图从表格中选择。 保持选中重置所有检查点(默认设置),以恢复源偏移量、算子状态和表数据,从而使流水线从回退点重新开始处理。 清除它只用于恢复表数据,不重新处理,比如当你想恢复表内容但不想重新处理受影响的数据时。 对于表及其上游,此设置必须保持相同;对于读取 Kafka 或 Auto Loader 等外部源的流,则不能清除此设置。 单击 “下一步”
  4. 查看回退点、检查点设置和受影响的数据集,然后点击“回退”。

启动流水线以重放数据。

你可以反复倒带。 每次倒回都会替换上一次的倒回,因此,如果回放失败,你可以通过倒回到其他位置来恢复。

倒带后

回放失败会导致流水线回退到先前位置,但处于停止状态。 修复代码或源数据,然后重新启动流水线以重试,或者回退到其他时间点。 流水线不会自行回滚,错误通过标准流水线诊断和事件日志显现。

只有当前管道定义与已恢复的状态兼容时,重放才会成功;倒带不会放宽 Structured Streaming 的兼容性规则。 关于哪些更改是兼容的,请参见 结构化流查询中的变更类型。 具体化视图遵循批处理语义并容忍更广泛的模式变更,但如果依赖不兼容,仍然会失败。

回退操作如果在过程中失败,可能会导致管道仅部分回退。 您有两个选项:

  1. 再倒回同一点或另一个点,管道就会汇聚到那个点。
  2. 要在回退未完成的情况下强制流水线启动正常更新,请将 pipelines.allowUpdateAfterIncompleteRewind 设置为 true,然后重启流水线。

你最多可以回退多远

回退点将保留 7 天。 在该时间窗口内,如果回退所需的数据已被移除,则回退将失败。 在依赖 Rewind 之前,请先检查以下几点:

  • VACUUM,或对您的表缩短 delta.deletedFileRetentionDuration。 请参阅处理表历史记录
  • 源数据保留时长短于你想要回溯的时间窗口,例如某个只保留一天数据的 Kafka 主题。

包含多个数据源或较长依赖链的管道需要更长的保留期,因为每个表和检查点都必须能够回溯到同一个一致的时间点。

局限性

  • Rewind 无法将流水线恢复到执行完全刷新之前的时间点。
  • 回退点会保留 7 天;如果回退到该时间点所需的表历史记录或源数据已被移除,则回退将失败,例如由于 VACUUM 或源数据保留期过短。 看看 你能倒回多久
  • 不支持实时模式。
  • Kinesis、Pulsar、Google Pub/Sub 以及使用 DSv2 或 Python 数据源 API 构建的自定义源码不被支持作为源码。
  • 不支持外部汇和自定义汇,包括定义为 create_sink()的汇。 请参阅 在管道中使用接收器
  • 使用行筛选或列掩码的流式表无法回溯。 请参阅 手动应用行筛选器和列掩码
  • 有状态回退需要 RocksDB 状态存储,而管道默认使用该存储。 对于配置了不同状态存储的流,回退会失败。
  • 某些 AUTO CDC 流式表需要先刷新,然后才能回退。 当你请求回退时,管道会通知你。
  • 实体化视图在回退后可能会执行完全重新计算,而不是进行增量刷新。 请参阅具体化视图的增量刷新
  • Rewind 仅适用于一个管道,不会与外部读取程序或其他读取同一表的管道进行协调。
  • Rewind 不会恢复流水线代码、流水线配置或 Unity 目录对象元数据,如标签和授权。

其他资源