返回
RSS Databricks Blog AI 逐段翻译 发布 2026-09-15 05:04

Databricks 发布 Spark 结构化流状态按需重分区

DataHot 速览

Databricks 宣布在 Databricks Runtime 18 及以上版本中公开预览“按需状态重分区”功能,允许有状态 Apache Spark Structured Streaming 查询调整分区数而保留 checkpoint 状态。此前,分区数在流启动时写入 checkpoint,修改 spark.sql.shuffle.partitions 并重启查询不会生效;若想应用新分区数,只能放弃现有 checkpoint 并从头开始,这会丢失累积状态。该功能适用于聚合、流-流 join、去重、会话化和 transformWithState 等有状态查询。早期用户 Coveo 表示,其相关 Amazon S3 API 成本降低了 40%,此前扩容决策常在过度配置与重建 checkpoint 之间权衡。

为什么值得关注:对有状态流处理生产用户而言,这解决了长期存在的分区数无法在线调整、调整即丢状态的问题,并给出成本与性能优化路径;Coveo 的 40% 成本下降提供了可量化的运营参考。

本文目录 7 节
  1. 深入底层:为什么状态分区会被锁死?
  2. 入门所需条件
  3. 更改分区数量
  4. 监控重新分区操作
  5. 示例:将查询缩容
  6. 何时使用状态重新分区
  7. 结论

译文

AI 逐段翻译

任何在生产环境中运行有状态 Apache Spark™ Structured Streaming 查询的人最终都会撞上同一堵令人不适的墙。

你在几个月前启动了该查询。当时数据量不大,所以你接受了 200 个 shuffle 分区的默认值就继续了。流水线运行得很顺畅。后来业务增长了,流量翻了三倍,状态存储也膨胀了。突然之间,那 200 个分区不再合适了。有些分区出现倾斜并运行过热,集群不堪重负,每个微批次所花的时间都超出了应有水平。

于是你做了最自然的举动:调高 spark.sql.shuffle.partitions 并重启查询。结果什么都没变。

查询悄悄地忽略了你设置的新值,因为分区数在你首次启动流时就已经固化到检查点里了。历来,应用新数值的唯一办法就是放弃现有检查点并重新开始,而对于有状态查询来说,这意味着丢失你一直精心维护的所有累积状态。对于一个追踪数百万账户的欺诈模型,或者一个持有数天窗口的会话化作业来说,“重新开始”绝不是任何人在生产事故复盘会上想说的话。

按需状态重新分区(公共预览版)已在 Databricks Runtime 18 及更高版本中提供,它移除了这堵墙。你现在可以为有状态流式查询调整分区数量,同时保持检查点状态完好无损。

这适用于任何有状态流式查询,无论你运行的是聚合、流-流连接、去重、会话化还是 transformWithState,也适用于任何工作负载,从欺诈检测到实时监控。

对于像 Coveo 这样的早期采用者来说,按需调整其流式基础设施规模的能力立即转化为了显著的运营成本节约。

在 Coveo,我们运行大规模的有状态流式流水线,数据量会随时间大幅波动。借助 Databricks 和状态重新分区能力,我们将相关的 Amazon S3 API 成本削减了 40%。以前,每一次扩缩容决策都迫使我们在两难中取舍:要么过度配置,要么从新的检查点重建,这会让存储 API 成本几乎与计算成本相当。现在我们随着需求变化自由扩缩,而不会中断现有状态或触发代价高昂的检查点迁移。”——Alexis Chicoine,Coveo 高级软件开发工程师

深入底层:为什么状态分区会被锁死?

要理解 Coveo 的成果为何代表 Structured Streaming 的一次有意义的飞跃,我们必须先看看分区数最初为什么会被冻结。

有状态流式查询将其状态保存在状态存储中,而该状态在物理上是分区的。流中的每个键——用户 ID、账号、窗口——都会哈希到特定分区,每个分区的数据都存储在检查点内各自独立的 RocksDB 实例中。分区数量定义了整个状态存储在磁盘上的布局。

如果你只是在两次重启之间更改分区数,哈希就不再对齐了。一个之前位于某个分区(比如 47 号分区)的键,现在可能哈希到另一个分区(12 号分区),但它累积的状态仍然留在原分区的文件里。实际上,查询会丢失对自身记忆的追踪。为了防止正是这种静默损坏,Structured Streaming 在创建检查点时锁定了分区数,并忽略之后对 spark.sql.shuffle.partitions 的任何更改。

安全,但不灵活。你付出的两项代价是:

  1. 你无法调优。 如果 200 个分区被证明是错误的选择,那么在整个检查点生命周期内你都被困在其中。
  2. 你无法随工作负载扩缩容。 随着数据量增长或缩减,你的分区数无法跟上步伐。

按需状态重新分区解决了这两个问题,它做了旧设计拒绝做的那一件事,但以安全的方式完成——通过物理地重新分布状态以匹配新的分区数。

入门所需条件

要求很简短:

  • Databricks Runtime 18 或更高版本。
  • RocksDB 状态存储提供程序。 在 DBR 17.3 及更高版本中,RocksDB 是默认选项,在这些版本中创建的新查询将使用它,除非显式更改。如果你想确认或显式设置它,请参见 在 Databricks 上配置 RocksDB 状态存储

这就是完整的先决条件列表。如果你使用的是 DBR 18 及默认状态存储,那么你已经拥有了所需的一切。

更改分区数量

该机制很简单,而且复用了每个流式开发者都已经熟悉的模式:停止、重新配置、重启。

不再设置 spark.sql.shuffle.partitions,而是设置一个专用配置 spark.sql.streaming.stateStore.partitions,然后重启查询:

关键细节在于这个新配置本身。对于有状态查询,spark.sql.streaming.stateStore.partitions 优先于 spark.sql.shuffle.partitions。这正是让该更改“生效”而旧方法做不到的原因。

当查询重启时,它不会立即恢复正常处理。首先,它会完成最后计划的微批次(如果还有待处理的)。然后它执行一次性重新分区操作:将状态数据物理地重新分布到新的分区数量上,将键重新哈希到它们正确的新归属,确保没有任何丢失或错位。一旦重新分布完成,查询就照常恢复处理,此时使用你请求的分区数。

这个重新分区步骤是该功能的核心。这就是“我们改了一个数字”与“我们安全地把你的状态迁移到了新布局”之间的区别。

监控重新分区操作

由于重新分区是一个实际的操作,其运行时间与状态量成正比,你会希望获得对它的可见性。Structured Streaming 通过其标准进度报告来呈现这一点。

在下一个微批次完成后,StreamingQueryProgress 事件会包含重新分区操作的持续时间。请在事件的 durationMs 指标中查看 controlBatch.REPARTITION字段,它报告以毫秒为单位的分区调整持续时间。

状态占用越大意味着分区调整耗时越长,但我们预计大多数工作负载只需几秒钟。因此,在大型作业上,值得捕获此指标以了解其持续时间。有关读取这些事件的更多信息,请参阅 在 Databricks 上监控 Structured Streaming 查询

示例:将查询缩容

让我们通过一个简单的聚合来具体说明,即按 id 对事件进行滚动窗口计数。我们将以默认的 200 个分区启动它,判定这超出了此工作负载所需,然后将其缩容到 100。

首先,是查询当前运行的状态,使用默认分区数:

现在,我们观察此流一段时间后得出结论:200 个分区属于过度配置。我们正在为我们不需要的并行度支付协调开销。我们停止查询,设置新的分区数,并使用相同的选项和相同的检查点重新启动它:

当重新启动的查询启动时,它会完成最后计划的一个微批次(如果仍有待处理的微批次),执行分区调整以将状态从 200 个分区重新分布到 100 个分区,然后继续计数,同时完整保留每个窗口和每个运行总计。同样的过程也适用于反向操作:要在更重的负载下扩容,只需设置一个更大的数字即可。

同样的方法也适用于 Spark 声明式管道(SDP)。有关完整演练,请参阅文档中的 SDP 示例

何时使用状态重新分区

按需状态重新分区是一种调优和扩缩工具,而非日常操作。它在以下几种关键情况下很有价值:

  • 上线后的规模调整。你在第一天就以默认的 200 个分区启动了管道,因为当时流很小,不值得进行微调。六个月后,这个数字已经固化在了一个你无法承受丢失的检查点中,而且它已经不够用了。例如:一个最初在单个试点区域启动的欺诈评分流现在已覆盖所有市场,而 200 个分区使得每个分区都持有过多的状态。借助按需重新分区,你可以增加分区数,以匹配你现在承载的数据量,而不会丢失现有检查点。
  • 工作负载变化。你按峰值流量规划了流的大小。例如:一个广告竞价管道在白天运行繁忙,在夜间则变得安静,因此针对白天峰值调优的值会导致凌晨 3 点大多数分区处于空闲状态。借助按需重新分区,你可以在进入繁忙时段时扩容,在繁忙时段过去后缩容,从而使分区跟随实际负载而不是最坏情况。
  • 回填历史数据:回填和稳态处理需要不同的分区数,而以前,你必须在检查点的整个生命周期内选择其中一个。例如:重新处理两年的历史数据需要高分区数来分散工作并快速完成,但一旦恢复到稳态流量,同样的分区数就会造成浪费。按需重新分区让你能够为回填扩容,并在赶上进度后缩容到稳态大小,而这一切都不会丢失检查点和状态。
  • 性能调优。分区数会影响并行度、状态大小和洗牌开销,而最优值很难预测。例如:你可能认为 200 太小,400 会降低微批次延迟,但过去的测试需要重建状态并重新处理数据,浪费资源。按需重新分区让你可以针对实时检查点调整分区数,并监控 controlBatch.REPARTITION 和微批次持续时间,从而基于测量结果而不是猜测来做出决定。

由于每次更改都需要停止并重新启动,并伴随一次性的重新分区暂停,请将其视为一项深思熟虑的维护操作。请将规模调整计划在可接受短暂处理暂停的时间窗口内进行,并观察 controlBatch.REPARTITION以确认其耗时,并让查询恢复正常的节奏。

结论

多年来,有状态流式查询的分区数是一个你在最开始就一次性做出的决定,之后要么再也不重新考虑,要么付出高昂代价从头重建状态。按需状态重新分区消除了这些限制。安全地在新分区数之间重新分布状态,将一项仅在启动时做出的决定转变为一项可在工作负载需要时随时重新考虑的决定。

其结果正是长期运行流的运维人员所期望的:能够根据其扩缩需求自由地对查询进行规模调整,只需停止、更改配置并重新启动,而不会丢失其状态。

按需状态重新分区在 Databricks Runtime 18 及更高版本中可用,使用 RocksDB 状态存储提供程序。有关完整参考,请参阅 有状态流式查询的按需状态重新分区

这篇内容对你有用吗?

反馈只用于改善内容筛选,不等同于收藏

分享这条资讯
分享海报
保存图片
iOS 也可以长按图片保存