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

在 Azure Databricks 上以计划的 Lakeflow 作业的形式运行生产结构化流式处理工作负荷。 请参阅 Lakeflow Jobs

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

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

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

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

注意

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

注意

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

对于无服务器计算上的连续流式处理,请在连续模式下使用触发与连续管道模式。

请参阅 流式处理限制

降低运营流传输的延迟

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

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

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

  • 输出模式:在查询运算符和接收器支持的情况下,使用更新模式。 更新模式在每次触发后都会输出更新后的行,并持续更新这些行,直到水位线过期,因此请确保下游接收端具有幂等性,以处理这些更新结果。 对于更新模式不支持的工作负载(例如流-流联接),或者在可以丢弃延迟到达的数据时,请使用附加模式。 为了低延迟,不要用完整模式。 请参阅为结构化流式处理选择输出模式

  • 触发器:使用 processingTime 带有 0 区间的触发器,当上一个微批次结束且有新数据可用时,立即开始下一个微批次。 这提供了最低的微批处理延迟,但增加了云存储API成本。 不要用 AvailableNowOnceContinuous 用于运营工作负载。 请参阅配置结构化流式处理触发器间隔

  • 水印:将水印时长设置得足够长,以涵盖工作负载不可丢弃的延迟到达数据。 水印控制查询接受事件时间错序数据的时间长度,之后会丢弃并驱逐状态,因此水印过短会无声地丢弃有效的迟到记录。 在此限制下,较短的水印降低延迟并保留更少状态,较长的水印容忍更多延迟数据,但代价是延迟和状态。 将延迟 SLA 乘以一个较小的倍数(例如 2 倍),可作为调优的合理起点。 请参阅应用水印来控制数据处理阈值

  • 源和接收端:从低延迟源读取数据,例如消息总线(Apache Kafka、Amazon Kinesis 或 Apache Pulsar),或者读取来自 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 服务跟踪流式处理查询,这会影响积压指标和作业通知。

在以下情况下使用 awaitTermination()

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