大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载导读当 Flink 作业处于严重背压backpressure时对齐 Checkpoint 的端到端耗时往往被 Barrier 穿越数据通道的时间所主导导致 Checkpoint 周期异常拉长、作业故障恢复风险上升。本文基于 Apache Flink 官方运维文档《Checkpointing under backpressure》系统讲解背压场景下 Checkpoint 变慢的根因与三种应对思路并重点深入两个可落地的优化方案缓冲区 DebloatingBuffer Debloating与非对齐 CheckpointUnaligned 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 监控指南。当这种情况发生并成为一个问题时有三种方法可以解决消除背压源头通过优化 Flink 作业逻辑、调整 Flink 或 JVM 参数抑或是对作业进行扩容rescaling来直接消除背压减少 In-flight 数据量降低 Flink 作业中缓冲在途in-flight数据的数据量启用非对齐 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.javaJavaStreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); // 启用非对齐 Checkpoint env.getCheckpointConfig().enableUnalignedCheckpoints();Scalaval env StreamExecutionEnvironment.getExecutionEnvironment() // 启用非对齐 Checkpoint env.getCheckpointConfig.enableUnalignedCheckpoints()Pythonenv 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同样已标记为弃用。限制一并发 CheckpointFlink 当前并不支持并发的非对齐 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 的耗时高于预期或产生波动典型场景包括一次性触发大量定时器timer例如在窗口操作中大量定时器同时触发时当前记录的处理时间会被显著拉长单条输入记录需要等待多个网络缓冲区系统在等待网络缓冲区可用时被阻塞例如序列化一条超出单个网络缓冲区容量的大记录时在flatMap操作中一条输入记录产生大量输出记录时背压会阻塞非对齐 Checkpoint直到处理这条输入记录所需的全部网络缓冲区都可用为止。任何其他单条记录处理耗时较长的场景也都可能触发该问题。限制四某些数据分布模式无法被 Checkpoint 覆盖有一部分包含特定属性的连接无法与 Channel 中的数据一样保存在 Checkpoint 中。为了保留这些特性并且确保没有状态冲突或非预期的行为非对齐 Checkpoint 对于这些类型的连接是禁用的所有其他的交换exchange仍然执行非对齐 Checkpoint。点对点连接Pointwise connections我们目前没有任何对于点对点连接中有关数据有序性的强保证。然而由于数据已经被以前置的 Source 或是 KeyBy 相同的方式隐式组织一些用户会依靠这种特性在提供有序性保证的同时将计算敏感型的任务划分为更小的块。只要并行度不变非对齐 CheckpointUC将会保留这些特性但是如果加上 UC 的扩缩容这些特性将会被改变。如上图所示的任务中如果我们想将并行度从 p2 扩容到 p3那么需要根据 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 数据内部的损坏导致无法恢复的情况否则不要设置该属性。另外请注意两点只有在重新部署作业后该属性才会生效这就意味着只有启用了 externalized checkpoint外部化 Checkpoint 时此操作才有意义——你需要能够从某个已完成的 Checkpoint 恢复该配置在源码中的正式 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/OCheckpoint 变大对齐超时降级想兼顾对齐的稳定性与非对齐的兜底背压持续时实际仍以非对齐为主实践要点回顾非对齐 Checkpoint 仅在EXACTLY_ONCE 最大并发 Checkpoint 为 1时可用缓冲区 Debloating 的target参数直接决定了 In-flight 数据的目标量级period与samples控制调整的敏捷度与稳定性非对齐 Checkpoint 与 Watermark、点对点/广播连接、长记录处理存在交互限制生产环境应先在测试作业上验证行为差异出现 In-flight 数据损坏等极端情况时execution.state-recovery.without-channel-state.checkpoint-id是最后的恢复手段但会导致数据丢失务必慎用。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐黑苹果配置革命OpCore Simplify让你15分钟搞定专业级EFI黑苹果配置革命OpCore Simplify让你15分钟搞定专业级EFI 还在为黑苹果配置头疼吗面对复杂的ACPI补丁、内核扩展、设备属性你是不是感觉像是开发工具CLIFlink 大状态与 Checkpoint 调优实战指南从 RocksDB 内存到 Task 本地恢复Flink 大状态与 Checkpoint 调优实战指南从 RocksDB 内存到 Task 本地恢复 本文是 Flink 流处理运维场景中针对 大状态La大数据流处理批处理数据工程torchtitan-npu Checkpoint 完整使用指南DCP 断点续训、Hugging Face 权重加载与 seed checkpoint 实战torchtitan npu Checkpoint 完整使用指南DCP 断点续训、Hugging Face 权重加载与 seed checkpoint 实战人工智能大模型分布式训练预训练模型优化Ascend创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
