结构化流式处理的生产注意事项

在 Azure Databricks 上将生产环境的 Structured Streaming 工作负载作为已调度的 Lakeflow 作业运行。 请参阅 Lakeflow Jobs

Databricks 建议始终配置以下内容:

  • 从返回结果的笔记本中删除不必要的代码,例如 displaycount
  • 不要使用全用途计算运行结构化流式处理工作负荷。 始终使用 jobs compute 将流作为 Lakeflow 作业进行调度。
  • 使用 Continuous 模式调度 Lakeflow 作业。 此处指的是 Azure Databricks 的作业调度功能,而非 Structured Streaming 的 触发间隔
  • 请勿为 Structured Streaming 作业的计算资源启用自动缩放。

某些工作负载会受益于以下功能:

Databricks 引入了 Lakeflow 管道,以减少管理结构化流式处理工作负荷的生产基础结构的复杂性。 Databricks 建议对于新的结构化流管道使用 Lakeflow 管道。 请参阅 Spark Declarative Pipelines

注意

计算自动缩放在缩减结构化流式处理工作负载的群集大小方面存在限制。 Databricks 建议针对流式工作负载,在 Lakeflow 上使用具有增强型自动扩缩功能的 Spark 声明式管道。 请参阅 使用自动缩放优化 Lakeflow 管道群集利用率

:::note 无服务器计算

在无服务器计算中,仅 Trigger.AvailableNow() 受支持且 Trigger.Once() 受支持。 Databricks 建议使用Trigger.AvailableNow()

若要在无服务器计算环境中进行连续流处理,请在连续模式下使用触发式与连续管道模式

请参阅 流式处理限制

:::

降低运营流传输的延迟

操作型流式工作负载近实时地摄取、转换数据并据此采取行动。 常见的例子包括欺诈检测、异常检测、个性化以及实时监控和警报,这些延迟处理直接影响业务成果。 这些工作负载的低延迟通常意味着数十到数百毫秒,尽管许多团队会在秒级范围内设定服务水平协议(SLA),以考虑高百分位时的变异性。

为了获得最低端到端延迟,可以使用实时模式,该模式在尾端延迟低于1秒,常见情况下约为300毫秒。 参见 实时模式概念

当实时模式不适合您的工作负载时,以下最佳实践可降低微批次结构化流的延迟:

  • 输出模式:在查询运算符和接收器支持的情况下,使用更新模式。 更新模式在每次触发后都会输出更新后的行,并持续更新这些行,直到水位线过期,因此请确保下游接收端具有幂等性,以处理这些更新结果。 对于更新模式不支持的工作负载(例如流-流联接),或者在可以丢弃延迟到达的数据时,请使用附加模式。 为了低延迟,不要用完整模式。 请参阅为结构化流式处理选择输出模式
  • 触发器:使用 processingTime 带有 0 区间的触发器,当上一个微批次结束且有新数据可用时,立即开始下一个微批次。 这提供了最低的微批处理延迟,但增加了云存储API成本。 不要用 AvailableNowOnceContinuous 用于运营工作负载。 请参阅配置结构化流式处理触发器间隔
  • 水印:将水印时长设置得足够长,以涵盖工作负载不可丢弃的延迟到达数据。 水印控制查询接受事件时间错序数据的时间长度,之后会丢弃并驱逐状态,因此水印过短会无声地丢弃有效的迟到记录。 在此限制下,较短的水印降低延迟并保留更少状态,较长的水印容忍更多延迟数据,但代价是延迟和状态。 将延迟 SLA 乘以一个较小的倍数(例如 2 倍),可作为调优的合理起点。 请参阅应用水印来控制数据处理阈值
  • 源和汇:从低延迟源读取数据,例如消息总线(Apache Kafka、Amazon Kinesis、Apache Pulsar 或 Google Cloud Pub/Sub),或者从 Delta Lake 和 Apache Iceberg 表的变更数据流中读取数据。 写入到低延迟、高吞吐量的接收器,如消息总线、操作数据库或 foreach 接收器。 将接收端操作设计为幂等,以便下游消费者能够处理重复数据和迟到数据。
  • 状态与检查点:对于有状态查询,使用RocksDB状态存储,该存储对变更日志检查点和异步状态检查点均为必需。 启用变更日志检查点,只持久化增量状态变化。 当状态检查点成为批次持续时间的瓶颈时,在了解异步状态检查点在故障恢复和集群扩缩容方面的注意事项后,启用异步状态检查点,以使检查点写入与下一个微批次并行进行。 在持久的云存储中,给每个查询单独设置检查点目录。 参见在 Azure Databricks 上配置 RocksDB 状态存储有状态查询的异步状态检查点结构化流检查点
  • 偏移管理:为了减少连续流偏移检查点带来的延迟,启用异步进度追踪,更新偏移和提交日志而不阻断数据处理。 它与 OnceAvailableNow 触发器不兼容。 请参阅 异步进度跟踪
  • 存储跳转:尽可能在单个流式管道内完成计算。 将逻辑拆分到多个作业或管道会增加存储跳数,从而增加延迟。

设计流式处理工作负载来应对失败

Databricks 建议始终将流式处理作业配置为在失败时自动重启。 某些功能(包括架构演变)要求结构化流式处理工作负载自动重试。 请参阅如何配置结构化流式处理作业以在失败时重启流式查询

有些操作(例如 foreachBatch)提供至少一次(而不是恰好一次)保证。 对于这些操作,请确保处理管道是幂等的。 请参阅使用 foreachBatch 将内容写入到任意数据接收器

注意

当查询重启时,将会处理在之前运行中计划的微批处理。 如果您的作业由于内存不足错误导致失败,或者您因微批次过大而手动取消作业,则可能需要升级计算资源,以便成功处理微批次。

如果在运行之间更改了配置,这些配置将应用于计划的第一个新批处理。 请参阅在结构化流式处理查询发生更改后恢复

作业重试时

可以将多个任务安排为Azure Databricks作业的一部分。 使用连续触发器配置作业时,无法设置任务之间的依赖项。

可选择使用以下方法之一在单个作业中计划多个流:

  • 多任务:定义一个具有多个任务的作业,这些任务会使用连续触发器运行流式处理工作负载。
  • 多查询:在单个任务的源代码中定义多个流式处理查询。

还可以组合使用这些策略。 下表比较了这些方法。

策略 多个任务 多个查询
如何共享计算? Databricks 建议为每个流式处理任务部署适当大小的计算资源。 可以选择跨任务共享计算。 所有查询共享相同的计算。 可以选择将查询分配给 调度池
如何处理重试? 在作业重试之前,所有任务都必须失败。 如果任何查询失败,任务将会重试。

有关处理多个任务或查询的更多详细信息,请参阅 在同一群集上运行多个结构化流式处理查询

将结构化流式处理作业配置为在失败时重启流式处理查询

Databricks 建议将所有流式工作负载配置为使用连续触发器。 请参阅连续运行作业

默认情况下,连续触发器具有以下行为:

  • 防止作业同时多次运行。
  • 在上一次运行失败时启动新的运行。
  • 使用指数退避进行重试。

Databricks 建议在计划工作流时始终使用作业计算而不是通用计算。 在作业失败并重试时,将会部署新的计算资源。

注意

Databricks 建议不要使用 streamingQuery.awaitTermination()spark.streams.awaitAnyTermination()。 请参阅 何时使用 awaitTermination()

何时使用 awaitTermination()

streamingQuery.awaitTermination()spark.streams.awaitAnyTermination() 阻止当前线程,直到流式查询终止。 是否使用这些函数取决于执行环境。

在 Lakeflow 作业中,请勿使用 streamingQuery.awaitTermination()spark.streams.awaitAnyTermination()。 这些函数并非必需,因为当流式查询处于活动状态时,Jobs 服务会自动阻止任务运行完成。 这两个函数都会阻止笔记本单元格完成执行,并阻止 Jobs 服务跟踪流式查询,从而干扰积压指标和任务通知。

在以下情况下使用 awaitTermination()

用例 行为
用于全用途计算的交互式笔记本 awaitTermination() 使单元格保持运行状态,使你能够观察查询状态,并确保笔记本输出中的故障浮出水面。
本地和开发环境 在本地运行 Spark 程序时,当主线程完成时,进程将退出。 调用 awaitTermination() 以使程序保持活动状态,直到流式处理查询完成或失败。
故障蔓延至驱动程序 如果没有 awaitTermination(),那么在非作业上下文中的流式查询失败可能不会传播到调用线程。 查询可能会以无提示方式失败,从而使故障更难检测和诊断。 调用 awaitTermination() 会再次引发驱动程序上的查询异常。