如何找出 Flink 背压的源头?
作者:Piotr Nowojski(@PiotrNowojski) · 原文发表于 2021 年 7 月 7 日 · 来源:Apache Flink 官方博客:How to identify the source of backpressure?。
版本说明:本文翻译原文的技术内容,并加入明确标注的核查补充。原文围绕 Flink 1.13 的指标和 Web UI 展开。2026-10-09 查看 Flink 2.3 背压监控文档时,文档仍列出三项子任务时间指标,并明确说明任务图对背压和忙碌指标取子任务最大值及相应颜色含义;该段没有说明空闲指标在任务图中的聚合方式。下文对原文描述的当时聚合行为作历史说明,不应直接推广到所有部署版本。具体界面、刷新频率与数据源行为应按实际版本核对。原文结尾的升级建议仅反映 2021 年背景。
过去几年,人们从不同角度讨论过背压。但在识别和分析背压来源这件事上,最近几个版本的 Flink 已经有了很大变化,尤其是 Flink 1.13 新增的指标与 Web UI 功能。本文解释这些变化,并进一步说明如何沿着作业找到背压源头。在此之前,先回答一个问题。

什么是背压?
Ufuk Celebi 的一篇旧文章对背压有很清楚的解释,虽然发表已久,基本概念依然适用。如果还不熟悉背压,作者建议先阅读它。原文还链接了一篇对 Flink 网络栈与背压机制的底层解释,供希望深入理解实现的人参考。核查说明:前一个 Ververica 旧链接目前重定向至其博客首页,本稿保留原引用,但不把重定向页面当成该历史文章正文。
从整体上看,当作业图中的某些算子处理记录的速度,赶不上收到记录的速度时,就会出现背压。运行这个慢算子的子任务,其输入缓冲区会逐渐填满。输入缓冲区满了以后,压力就传递到上游子任务的输出缓冲区。
当上游的输出缓冲区也填满,上游子任务便不得不降低记录处理速度,使之与下游瓶颈算子的处理速度一致。背压继续沿着数据流反方向传播,最终可以一路传到源算子。
如果负载与可用资源保持稳定,而且没有算子短时间内突发输出大量数据,例如窗口算子触发输出,那么输入、输出缓冲区一般只会处于两种状态之一:几乎为空,或者几乎全满。下游算子或子任务跟得上流入速度时,缓冲区为空;跟不上时,缓冲区则会填满。网络交换本身成为瓶颈时还有一个例外,见文末注释。
检查缓冲区使用率,正是 Nico Kruber 几年前介绍的背压检测与分析方法的基础。现在 Flink 提供了更直接的工具,但在介绍工具之前,还有两个问题值得考虑。
为什么需要关心背压?
背压是机器或算子负载过高的一个信号。背压累积会直接影响系统端到端延迟,因为记录在真正得到处理之前,需要在队列里等待更长时间。
其次,背压会使对齐检查点花费更长时间,而非对齐检查点则会变大。有关两种检查点的背景,可以参见 Flink 1.13 的背压下检查点文档。如果你正受到检查点屏障传播时间过长的困扰,处理背压很可能有助于解决问题。还有一种动机更直接:你可能只是想优化作业,降低运行成本。
不论是哪一种情况,处理问题都需要先意识到它存在,再定位并分析它。
为什么有时可以不在意背压?
坦率地说,并不是只要出现背压就必须处理。几乎从定义上讲,没有背压,通常意味着集群至少还有一些未充分利用的资源,资源供给也略高于需求。若目标是尽量减少闲置资源,恐怕难以完全避免背压。批处理尤其如此。
审核补充:是否需要消除背压,应结合延迟目标、检查点行为与吞吐判断。“存在背压”不是“系统故障”的同义词。
如何检测并追踪背压源头?
一种办法是直接查看指标。但从 Flink 1.13 起,通常不必一开始就深入到这一层:多数情况下,先看 Web UI 中的作业图就足够了。
原文的作业图示例中,不同任务采用不同颜色。颜色同时表达两个因素:该任务受到多少背压,以及它有多忙。空闲任务是蓝色,完全忙碌的任务是红色,完全受背压的任务是黑色。介于这些状态之间的任务,会呈现三种颜色的混合或不同深浅。
掌握这层含义以后,受背压的黑色任务便很容易辨认。沿数据流向下游看,位于这些任务下游、最忙的红色任务,很可能就是背压源头,也就是瓶颈。
点击某个任务,进入 BackPressure 选项卡,还能进一步检查该任务中每个子任务处于忙碌、背压或空闲状态的比例。当存在数据倾斜、各个子任务的利用率并不一致时,这尤其有用。
原文的子任务截图展示了一种情形:可以清楚地分出哪些子任务空闲、哪些子任务受到背压,而且其中没有子任务处于忙碌状态。通常,仅靠这些信息就能很快理解作业大致在发生什么。不过,仍有一些细节需要解释。
这些数字究竟表示什么?
这套监控机制的基础,是每个子任务计算并暴露的三个指标:
| 指标 | 含义 |
|---|---|
idleTimeMsPerSecond |
平均每秒处于空闲状态的毫秒数。 |
busyTimeMsPerSecond |
平均每秒处于忙碌状态的毫秒数。 |
backPressuredTimeMsPerSecond |
平均每秒处于背压状态的毫秒数。 |
除去舍入误差,这三项彼此互补,在同一个子任务内应当合计约为 1000 ms/s。从表达时间占比这一点上看,它们与 CPU 使用率指标有些相似。
还要注意,这些数值是最近几秒的短期平均值,并且包括子任务线程内发生的各种工作:算子、函数、定时器、检查点、记录序列化与反序列化、网络栈,以及其他 Flink 内部开销。一个忙于触发定时器并产生结果的 WindowOperator,会被记为忙碌或受到背压。如果某个函数在 CheckpointedFunction#snapshotState 调用中进行开销较大的工作,例如刷新内部缓冲区,它也会被记为忙碌。
但 busyTimeMsPerSecond 和 idleTimeMsPerSecond 看不到主子任务执行循环之外、独立线程中发生的工作。原文指出两个相关场景:
- 用户在算子内部手工创建的线程。原文不鼓励这种做法。
- 实现已经弃用的
SourceFunction接口的旧式数据源。在原文版本中,这类数据源的busyTimeMsPerSecond会显示NaN或N/A。参见当时的数据源文档。
审核补充:这里的“忙碌”不是 CPU 使用率。阻塞调用、序列化、状态处理等工作都可能影响该值;需要结合机器资源、线程栈或性能分析判断原因。旧接口的可用性和展示行为则需要按部署版本确认。
为了在 Web UI 中展示这些原始数字,还需要对任务中的全部子任务进行聚合,因为作业图上显示的是任务,而不是每个子任务。原文将当时 Web UI 的聚合概括为取子任务中的最大值。Flink 2.3 文档明确说明任务图对背压和忙碌指标取最大值,但没有说明空闲指标在任务图中的聚合方式。即使只看已明确的背压和忙碌指标,最大值也可能来自不同子任务,因此两者不一定合计为 100%。
例如,一个子任务有 60% 的时间受到背压,而另一个子任务有 60% 的时间忙碌;任务级视图就可能同时显示“背压 60%”和“忙碌 60%”。这不是同一个子任务的时间重叠,更不是指标必然出错,而是最大值来自不同子任务。
负载变化会掩盖什么?
前面说过,指标是在几秒窗口内测量并求平均的。分析负载变化的作业时,要一直记住这一点,例如包含周期性触发的 WindowOperator 的任务或子任务。
一个始终保持 50% 负载的子任务,与另一个每隔一秒在完全忙碌和完全空闲之间切换的子任务,可能都报告 busyTimeMsPerSecond = 500 ms/s。平均数相同,运行模式却不同。
而且,负载变化,特别是窗口触发,还可能让瓶颈移到作业图的另一个位置。原文案例里,SlidingWindowOperator 在累积记录期间是瓶颈;一旦它开始每 10 秒触发一次窗口,下游任务 SlidingWindowCheckMapper -> Sink: SlidingWindowCheckPrintSink 就变成瓶颈,SlidingWindowOperator 反而开始受到背压。
因为忙碌、背压与空闲指标是在几秒内求平均,这种细微切换不一定能立刻从界面看出来,需要结合上下文推断。原文所用的 Web UI 每 10 秒才更新一次状态,也使更高频的变化更难观察。这里的 10 秒是原文环境的行为描述,不是对所有当前版本的固定保证。
发现背压以后,可以做什么?
如何处理背压是一个复杂主题,值得单独写一篇文章,早先的博客也已讨论过一部分。从总体上说,有两条路径:
- 增加资源:更多机器、更快 CPU、更多 RAM、更好的网络,或使用 SSD 等。
- 提高现有资源利用效率:优化代码、调整配置、避免数据倾斜。
无论采用哪条路径,都应先分析背压为何产生:
- 确认背压确实存在。
- 定位引发背压的子任务或机器。
- 继续深入,找出具体哪段代码造成问题,以及究竟哪种资源不足。
改进后的背压监控和指标可以帮助完成前两步。要做第三步,代码性能分析可能是合适的方法。为了方便分析,从 Flink 1.13 起,Web UI 集成了火焰图。火焰图是一种常用的性能分析与可视化方法,作者鼓励读者尝试。
找到瓶颈所在以后,分析方式其实与分析其他非分布式应用类似:查看资源利用情况、连接性能分析器等。通常没有通用的捷径。你可以尝试扩容,但有时扩容并不容易,也不实际。
背压监控的这些改进让背压源头更容易找到,而火焰图能帮助分析某个子任务为什么会成为问题。两者结合,可以显著简化过去比较繁琐的 Flink 作业调试和性能分析。原文最后邀请读者“升级到 Flink 1.13.x 并试用这些功能”;这是 2021 年的版本建议,现今应使用适合自身环境、仍受支持的版本及其对应文档。
审核补充:性能分析器可能引入开销,增加机器或提高并行度可能增加资源成本。实施优化后,需要用实际吞吐、延迟、资源利用和检查点变化来验证结果;本文没有运行作业、采集性能数据或验证任何性能收益。
注释与署名
缓冲区状态的例外:在较少见的“网络交换本身就是作业瓶颈”场景中,下游任务的输入缓冲区可能为空,而上游的输出缓冲区却已满。因此,不能只凭某个红色节点或缓冲区占用,就把根因归到该算子的业务代码。
原文及网站内容 © 2024 Apache Software Foundation,按 Apache License 2.0 提供(完整许可证文本);原文作者 Piotr Nowojski。Apache Flink、Flink 及其标志为 Apache Software Foundation 的商标或注册商标。修改说明:本文为原文的中文全文翻译,包含术语整理及明确标注的核验补充;原图已替换为本次绘制的技术示意图,未冒充原文 Web UI 截图。











暂无评论内容