Flink 背压下的 Checkpoint 优化:缓冲区 Debloating 与非对齐 Checkpoint 完整实战指南
2026/9/23 4:45:37 网站建设 项目流程
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

导读

当 Flink 作业处于严重背压(backpressure)时,对齐 Checkpoint 的端到端耗时往往被 Barrier 穿越数据通道的时间所主导,导致 Checkpoint 周期异常拉长、作业故障恢复风险上升。本文基于 Apache Flink 官方运维文档《Checkpointing under backpressure》,系统讲解背压场景下 Checkpoint 变慢的根因与三种应对思路,并重点深入两个可落地的优化方案:缓冲区 Debloating(Buffer Debloating)非对齐 Checkpoint(Unaligned Checkpoint)。读完本文,你将掌握如何通过配置与 API 启用这两项能力、如何用aligned-checkpoint-timeout让 Checkpoint 从对齐平滑降级为非对齐,并理解非对齐 Checkpoint 在并发、Watermark、数据分布等方面的限制与故障恢复手段,能够在真实背压作业中做出正确的取舍。

背压下 Checkpoint 变慢的根因

通常情况下,对齐 Checkpoint 的耗时主要由 Checkpoint 过程中的同步阶段异步阶段两部分决定。然而,当 Flink 作业正运行在严重的背压下时,Checkpoint 端到端延迟的主要影响因子会发生转移——传递 Checkpoint Barrier 到所有算子/子任务所需的时间成为决定性因素。

原因在于:在背压场景下,数据通道中的缓冲(buffer)被大量待处理数据填满,Checkpoint Barrier 必须排队等待前面的数据被消费后才能继续前进。Barrier 的传播速度因此被拖慢,整个 Checkpoint 的对齐时间(alignment time)被显著拉长。

在 Flink 的 Web UI 或监控指标中,可以通过 Checkpoint 监控页面的 History Tab 观察到两个关键指标来确认该问题:

  • Alignment time(对齐时间):Barrier 等待其他输入通道 Barrier 到达的时间;
  • Start delay(启动延迟):从 Checkpoint 触发到实际开始执行的时间。

如果这两个指标异常偏高,通常意味着背压已经严重拖慢了 Checkpoint 的推进。关于 Checkpoint 的整体流程,可参考 有状态流处理与 Checkpoint 概念;关于指标的具体查看方式,可参考 Checkpoint 监控指南。

当这种情况发生并成为一个问题时,有三种方法可以解决:

  1. 消除背压源头:通过优化 Flink 作业逻辑、调整 Flink 或 JVM 参数,抑或是对作业进行扩容(rescaling)来直接消除背压;
  2. 减少 In-flight 数据量:降低 Flink 作业中缓冲在途(in-flight)数据的数据量;
  3. 启用非对齐 Checkpoint:让 Checkpoint Barrier 不必等待数据通道中的数据被消费。

这些选项并不是互斥的,可以组合使用。本文重点介绍后两个选项。

缓冲区 Debloating:自动削减 In-flight 数据

特性概述与启用方式

缓冲区 Debloating 是 Flink 1.14 引入的一项新工具,用于自动控制Flink 算子/子任务之间缓冲的 In-flight 数据量。启用方式是在flink-conf.yml中设置:

taskmanager.network.memory.buffer-debloat.enabled: true

该配置项在源码中定义于 TaskManagerOptions.java,类型为布尔型,默认值为false(即默认关闭)。启用后,系统会根据实测吞吐量自动调整 In-flight 数据量。

对两种 Checkpoint 的影响

此特性对对齐非对齐Checkpoint 都生效,且在这两种情况下都能缩短 Checkpointing 的时间,不过 Debloating 的效果对于对齐 Checkpoint 最明显——因为 In-flight 数据减少后,Barrier 穿越数据通道的时间也随之缩短。

当在非对齐 Checkpoint情况下使用缓冲区 Debloating 时,还有一个额外的好处:Checkpoint 大小会更小,恢复时间更快。这是因为非对齐 Checkpoint 会把 In-flight 数据作为 Checkpoint State 的一部分持久化,In-flight 数据越少,需要保存和恢复的数据也就越少。

工作机制与配套参数

从实现层面看,缓冲区 Debloating 的核心逻辑位于 BufferDebloatConfiguration.java,它会读取TaskManagerOptions中定义的一组相关配置。其基本思想是:周期性测量当前网络吞吐量,再结合目标消费时间,动态计算出合适的 buffer 大小,从而把 In-flight 数据控制在"目标时间内可被完全消费"的量级。

围绕该特性,TaskManagerOptions.java 中定义了一组配套参数,全部位于 Task Manager 网络内存(Network Memory)配置区段:

配置项类型默认值说明
taskmanager.network.memory.buffer-debloat.enabledBooleanfalse自动缓冲区 Debloating 功能的总开关
taskmanager.network.memory.buffer-debloat.targetDuration1 s缓冲的 In-flight 数据应被完全消费的目标总时间。该值会与实测吞吐量结合,用于调整 In-flight 数据量
taskmanager.network.memory.buffer-debloat.periodDuration200 ms重新计算 buffer 大小的最小间隔周期。值越小对负载波动的反应越快,但可能影响性能
taskmanager.network.memory.buffer-debloat.samplesInteger20用于计算新 buffer 大小所采用的最近样本数量
taskmanager.network.memory.buffer-debloat.threshold-percentagesInteger25新计算的 buffer 大小与旧值之间的最小百分比差异,只有超过该差异才会应用新值,可避免频繁的小幅来回调整

例如,一段典型的 Debloating 配置如下:

taskmanager.network.memory.buffer-debloat.enabled: true taskmanager.network.memory.buffer-debloat.target: 1 s taskmanager.network.memory.buffer-debloat.period: 200 ms taskmanager.network.memory.buffer-debloat.samples: 20 taskmanager.network.memory.buffer-debloat.threshold-percentages: 25

需要注意的是,即使启用了缓冲区 Debloating,你仍然可以继续使用手动调优方式来减少缓冲在 In-flight 数据的数据量。关于缓冲区 Debloating 更完整的工作原理与调优细节,可参考 网络内存调优指南,该指南与本文配合阅读效果最佳。

非对齐 Checkpoint:让 Barrier 越过缓冲区

原理:Checkpoint 时长与吞吐量解耦

从 Flink 1.11 开始,Checkpoint 可以是非对齐的。非对齐 Checkpoint 会把 In-flight 数据(例如存储在缓冲区中的数据)作为 Checkpoint State 的一部分保存,从而允许 Checkpoint Barrier 跨越这些缓冲区,不必等待缓冲中的数据被消费完毕。因此,Checkpoint 时长变得与当前吞吐量无关——因为 Checkpoint Barrier 实际上已经不再嵌入到数据流当中了。

这意味着:如果你的 Checkpoint 由于背压导致周期非常长,就应该考虑使用非对齐 Checkpoint。启用后,Checkpointing 时间基本上与端到端延迟解耦。

但需要特别留意一个代价:非对齐 Checkpointing 会增加状态存储的 I/O。因为原本只需要持久化算子状态的 Checkpoint,现在还必须把大量的 In-flight 数据一并写入状态存储。因此,当状态存储的 I/O 是整个 Checkpointing 过程中的真正瓶颈时,你不应当使用非对齐 Checkpointing——此时它反而可能让情况更糟。

启用方式一:编程 API

在代码中可以通过CheckpointConfig直接启用。以下三种语言的 API 等价,均调用enableUnalignedCheckpoints()(对应方法定义于 CheckpointConfig.java):

Java

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 启用非对齐 Checkpoint env.getCheckpointConfig().enableUnalignedCheckpoints();

Scala

val env = StreamExecutionEnvironment.getExecutionEnvironment() // 启用非对齐 Checkpoint env.getCheckpointConfig.enableUnalignedCheckpoints()

Python

env = StreamExecutionEnvironment.get_execution_environment() # 启用非对齐 Checkpoint env.get_checkpoint_config().enable_unaligned_checkpoints()

启用方式二:配置文件

或者在flink-conf.yml配置文件中增加配置:

execution.checkpointing.unaligned: true

需要说明的是,源码中该配置的正式 key 为execution.checkpointing.unaligned.enabled(默认值false),而execution.checkpointing.unaligned是它的弃用别名(deprecated key),二者在配置文件中均可生效。对应配置定义于 ExecutionCheckpointingOptions.java,其描述明确指出:

启用非对齐 Checkpoint 可大幅缩短背压下的 Checkpoint 时长。非对齐 Checkpoint 只有在一致性模式为 EXACTLY_ONCE 且最大并发 Checkpoint 数为 1 时才能启用。

这一点非常重要:如果你的作业使用了 AT_LEAST_ONCE 或设置了execution.checkpointing.max-concurrent-checkpoints > 1,启用非对齐 Checkpoint 会失败。

对齐 Checkpoint 的超时:平滑降级机制

启用非对齐 Checkpoint 后,你依然可以指定对齐 Checkpoint 的超时,让每个 Checkpoint 在启动时先以对齐方式执行,超时后再降级为非对齐。这为两种模式提供了一种平滑的过渡策略。

通过编程方式设置:

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.getCheckpointConfig().setAlignedCheckpointTimeout(Duration.ofSeconds(30));

或者在flink-conf.yml配置文件中配置:

execution.checkpointing.aligned-checkpoint-timeout: 30 s

行为语义如下(定义见 ExecutionCheckpointingOptions.java 与 CheckpointingOptions.java):

  • 在启动时,每个 Checkpoint 仍然是aligned checkpoint(对齐 Checkpoint)
  • 但当全局 Checkpoint 持续时间超过aligned-checkpoint-timeout时,如果对齐 Checkpoint 还没完成,Checkpoint 将会转换为 Unaligned Checkpoint
  • 如果该超时设置为0,则 Checkpoint 将始终以非对齐方式启动(这也是默认值0 s的含义);
  • 该配置的旧 key 为execution.checkpointing.alignment-timeout,同样已标记为弃用。

限制一:并发 Checkpoint

Flink 当前并不支持并发的非对齐 Checkpoint。不过,由于非对齐 Checkpoint 带来了更可预测、更短的 Checkpointing 时长,实际场景中可能也根本不需要并发的 Checkpoint。

此外,Savepoint 也不能与非对齐 Checkpoint 同时发生,因此在这种组合下 Savepoint 将会花费稍长的时间。

限制二:与 Watermark 的相互影响

非对齐 Checkpoint 在恢复的过程中改变了关于Watermark 的一个隐式保证

目前,Flink 确保了 Watermark 作为恢复的第一步(而不是将最近的 Watermark 存放在 Operator 中),以便支持扩缩容。在非对齐 Checkpoint 中,这意味着:当恢复时,Flink 会在恢复 In-flight 数据后再生成 Watermark

如果你的 Pipeline 中使用了对每条记录都应用最新的 Watermark 的算子,那么使用非对齐 Checkpoint 将相对于使用对齐 Checkpoint 产生不同的结果。如果你的 Operator 依赖于最新的 Watermark 始终可用,解决办法是将 Watermark 存放在 OperatorState 中。在这种情况下,Watermark 应该使用单键 group 存放在 UnionState中,以方便扩缩容。

限制三:长时间记录处理(Long-running record processing)的影响

尽管非对齐 Checkpoint 的 Barrier 能够越过队列中的所有其他记录,但Barrier 的处理仍然可能被当前正在处理的那条记录阻塞——Flink 无法中断单条输入记录的处理过程,非对齐 Checkpoint 必须等待当前记录被完整处理完毕后才能继续。这会导致 Checkpoint 的耗时高于预期或产生波动,典型场景包括:

  1. 一次性触发大量定时器(timer):例如在窗口操作中,大量定时器同时触发时,当前记录的处理时间会被显著拉长;
  2. 单条输入记录需要等待多个网络缓冲区:系统在等待网络缓冲区可用时被阻塞,例如:
    • 序列化一条超出单个网络缓冲区容量的大记录时;
    • flatMap操作中,一条输入记录产生大量输出记录时,背压会阻塞非对齐 Checkpoint,直到处理这条输入记录所需的全部网络缓冲区都可用为止。

任何其他"单条记录处理耗时较长"的场景也都可能触发该问题。

限制四:某些数据分布模式无法被 Checkpoint 覆盖

有一部分包含特定属性的连接无法与 Channel 中的数据一样保存在 Checkpoint 中。为了保留这些特性并且确保没有状态冲突或非预期的行为,非对齐 Checkpoint 对于这些类型的连接是禁用的;所有其他的交换(exchange)仍然执行非对齐 Checkpoint。

点对点连接(Pointwise connections)

我们目前没有任何对于点对点连接中有关数据有序性的强保证。然而,由于数据已经被以前置的 Source 或是 KeyBy 相同的方式隐式组织,一些用户会依靠这种特性在提供有序性保证的同时,将计算敏感型的任务划分为更小的块。

只要并行度不变,非对齐 Checkpoint(UC)将会保留这些特性;但是如果加上 UC 的扩缩容,这些特性将会被改变

如上图所示的任务中,如果我们想将并行度从 p=2 扩容到 p=3,那么需要根据 KeyGroup 将 KeyBy 的 Channel 中的数据划分到 3 个 Channel 中去。这很容易做到,通过使用 Operator 的 KeyGroup 范围和确定记录属于某个 Key(group) 的方法(不管实际使用的是什么方法)。但对于Forward 的 Channel,我们根本没有 KeyContext——Forward Channel 里也没有任何记录被分配了任何 KeyGroup,也无法计算它,因为无法保证 Key 仍然存在。

广播连接(Broadcast connections)

广播连接带来了另一个问题:无法保证所有 Channel 中的记录都以相同的速率被消费。这可能导致某些 Task 已经应用了与特定广播事件对应的状态变更,而其他任务则没有。

广播分区通常用于实现广播状态(Broadcast State),它应该跨所有 Operator 都相同。Flink 实现广播状态的方式是:仅 Checkpointing 有状态算子的 SubTask 0 中状态的单份副本;在恢复时,将该份副本发送给所有的 Operator。因此,可能会发生以下情况:某个算子将很快从它的 Checkpointed Channel 消费数据并应用修改来获得状态,而其他算子尚未完成恢复,从而导致状态不一致的风险。

Troubleshooting:恢复损坏的 In-flight 数据

非对齐 Checkpoint 把 In-flight 数据也纳入了 Checkpoint 状态,因此理论上存在 In-flight 数据损坏导致无法恢复的极端情况。

⚠️ 警告:以下描述的操作是最后采取的手段,因为它们将会导致数据的丢失。

为了防止 In-flight 数据损坏,或者由于其他原因导致作业应该在没有 In-flight 数据的情况下恢复,可以使用recover-without-channel-state.checkpoint-id相关属性。该属性需要指定一个Checkpoint Id,对于它来说 In-flight 中的数据将会被忽略。

除非已经持久化的 In-flight 数据内部的损坏导致无法恢复的情况,否则不要设置该属性。另外请注意两点:

  1. 只有在重新部署作业后该属性才会生效,这就意味着只有启用了 externalized checkpoint(外部化 Checkpoint) 时,此操作才有意义——你需要能够从某个已完成的 Checkpoint 恢复;
  2. 该配置在源码中的正式 key 为execution.state-recovery.without-channel-state.checkpoint-id(默认值-1,即不忽略任何 In-flight 数据),文档中使用的execution.checkpointing.recover-without-channel-state.checkpoint-id是它的弃用别名。对应定义见 StateRecoveryOptions.java。

例如,如果要忽略 Checkpoint ID 为123456的 Checkpoint 中的 In-flight 数据,配置如下:

execution.state-recovery.without-channel-state.checkpoint-id: 123456

完整的配置项说明可参考 Flink 配置总览。

总结与选型建议

在背压导致 Checkpoint 周期过长的场景下,可以按以下顺序组合运用各项手段:

手段适用场景主要代价
消除背压源头背压可优化(作业逻辑、参数、扩容)需要投入调优/扩容成本
缓冲区 Debloating希望自动控制 In-flight 数据量,对对齐/非对齐均有效需要关注吞吐量波动对 buffer 调整的影响
非对齐 Checkpoint背压难以消除、Checkpoint 严重超时增加状态存储 I/O,Checkpoint 变大
对齐超时降级想兼顾对齐的稳定性与非对齐的兜底背压持续时实际仍以非对齐为主

实践要点回顾:

  • 非对齐 Checkpoint 仅在EXACTLY_ONCE + 最大并发 Checkpoint 为 1时可用;
  • 缓冲区 Debloating 的target参数直接决定了 In-flight 数据的目标量级,periodsamples控制调整的敏捷度与稳定性;
  • 非对齐 Checkpoint 与 Watermark、点对点/广播连接、长记录处理存在交互限制,生产环境应先在测试作业上验证行为差异;
  • 出现 In-flight 数据损坏等极端情况时,execution.state-recovery.without-channel-state.checkpoint-id是最后的恢复手段,但会导致数据丢失,务必慎用。
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

相关推荐

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询