尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Apache Flink 批作业推测执行(Speculative Execution)完整指南:原理、配置调优与 Source/Sink 适配

Apache Flink 批作业推测执行(Speculative Execution)完整指南:原理、配置调优与 Source/Sink 适配 Apache Flink 批作业推测执行Speculative Execution完整指南原理、配置调优与 Source/Sink 适配【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink导读本文围绕 Apache Flink 批处理作业的**推测执行Speculative Execution**机制展开讲解其产生的背景、底层工作原理、启用方式、参数调优策略以及如何让自定义 Source / Sink 与推测执行正确协作。读完本文你将掌握如何用一行配置为 Flink 批作业开启推测执行以抵御坏节点导致的作业变慢如何通过slow-task-detector系列参数精准调优慢任务检测如何通过 Web UI 与专用指标验证推测执行的实际效果以及如何改造自定义Source/Sink以兼容多并发执行尝试Execution Attempt。核心参考文档为 speculative_execution.md。背景为什么需要推测执行在分布式批处理集群中个别节点TaskManager可能出现硬件问题、突发的 I/O 繁忙或 CPU 负载过高。这些问题节点本身并不一定导致任务失败却会让其上运行的任务执行速度显著慢于其他节点上的同类任务最终拖慢整个批作业的执行时间。由于作业不会失败传统的失败重试机制对此无能为力只能被动等待慢任务完成。推测执行正是为缓解这类问题而设计当检测到某个任务执行过慢时Flink 会在未被判定为问题节点的其他节点上为该慢任务启动新的执行尝试attempt。新尝试与旧尝试消费相同的输入数据、产出相同的结果旧尝试不受影响、继续运行。最先完成的尝试被采纳其输出对下游任务可见并可被消费其余尝试随后被取消。工作机制慢任务检测 节点屏蔽 调度重部署从源码结构看推测执行在 Flink 运行时flink-runtime中由三部分协作完成慢任务检测器Slow Task Detector负责周期性识别慢任务。接口定义在 SlowTaskDetector.java当前实现为基于执行时间的 ExecutionTimeBasedSlowTaskDetector.java。节点屏蔽Blocklist机制慢任务所在的节点会被标记为问题节点并进入屏蔽列表调度器不会再把新的推测尝试部署到被屏蔽的节点上相关工具类见 BlocklistUtils.java。调度器创建并部署新尝试为慢任务创建新的执行尝试并调度到未被屏蔽的节点。该逻辑由 AdaptiveBatchScheduler.java 配合SpeculativeExecutionHandler实现类为 DefaultSpeculativeExecutionHandler.java另有用于关闭场景的 DummySpeculativeExecutionHandler.java完成。使用方式重要前提适用范围警告Flink 不支持对 DataSet 作业启用推测执行因为 DataSet API 将在不久后被废弃。DataStream API 是目前推荐的编写 Flink 批作业的低层 API。推测执行是面向**批作业Batch**的能力因此请确保作业基于 DataStream APIBatch 执行模式编写。启用推测执行只需在flink-conf.yaml或通过作业提交参数设置一个配置项execution.batch.speculative.enabled: true默认值为false参见 batch_execution_configuration.html。注意目前只有Adaptive Batch Scheduler自适应批调度器支持推测执行。Flink 批作业默认使用该调度器除非你显式配置了其他调度器。关于该调度器的更多说明见 elastic_scaling.md。调度器相关调优参数以下两个参数用于控制推测执行对调度的影响配置项默认值类型说明execution.batch.speculative.max-concurrent-executions2Integer每个算子可并发执行的最大执行尝试数量包含原始尝试和推测尝试。例如设置为 2意味着除原始尝试外最多再启动 1 个推测尝试。execution.batch.speculative.block-slow-node-duration1 minDuration被检测出的慢节点问题节点将被屏蔽Block多长时间。屏蔽期间调度器不会把新的推测尝试部署到该节点。慢任务检测器相关调优参数当前推测执行使用基于执行时间的慢任务检测器。以下参数控制检测的灵敏度与准确度配置项默认值类型说明slow-task-detector.check-interval1 sDuration慢任务检查周期即检测器每隔多久执行一次检测。slow-task-detector.execution-time.baseline-ratio0.75Double计算基线所需的已完成执行比例阈值 R。slow-task-detector.execution-time.baseline-multiplier1.5Double计算基线的放大倍数 M。slow-task-detector.execution-time.baseline-lower-bound1 minDuration慢任务检测基线Baseline的下限避免在作业刚启动、样本不足时把正常任务误判为慢任务。完整参数描述参见 slow_task_detector_configuration.html。基线Baseline的计算算法检测器会周期性统计所有**已完成finished**的执行。设算子并行度为 N配置比例为 R默认 0.75当已完成执行的比例达到N * R时取前N * R个已完成任务的执行时间中位数 T。基线 T × M其中 M 为slow-task-detector.execution-time.baseline-multiplier默认 1.5。当前仍在运行、且执行时间超过基线的任务即被判定为慢任务。数据倾斜Data Skew下的加权优化执行时间会按执行顶点Execution Vertex的输入数据量进行加权。因此当出现数据倾斜时输入数据量差异大但算力接近的执行不会被误判为慢任务从而避免启动不必要的推测尝试、浪费资源。警告如果算子节点是 Source或者使用了Hybrid Shuffle模式上述执行时间按输入数据量加权的优化不会生效因为此时无法获知输入数据量。让自定义 Source 适配推测执行当作业使用自定义 Source且该 Source 使用了自定义 SourceEvent 时需要让该 Source 的 SplitEnumerator 实现 SupportsHandleExecutionAttemptSourceEvent 接口public interface SupportsHandleExecutionAttemptSourceEvent { void handleSourceEvent(int subtaskId, int attemptNumber, SourceEvent sourceEvent); }该接口是SplitEnumerator的装饰性接口允许其处理来自特定执行尝试的SourceEvent见 SupportsHandleExecutionAttemptSourceEvent.java 的源码注释。这意味着SplitEnumerator必须能够感知到发送事件的到底是哪个尝试attemptNumber。否则当 JobManager 收到来自任务的 Source 事件时会发生异常导致作业失败。其他类型的 Source 无需任何额外改动即可配合推测执行包括SourceFunction 类型的 SourceInputFormat 类型的 Source新的 Source API 类型 Source。Apache Flink 官方提供的所有 Source Connector 都可以直接配合推测执行工作。让自定义 Sink 适配推测执行出于兼容性考虑Sink 默认不参与推测执行除非它实现了 SupportsConcurrentExecutionAttempts 接口public interface SupportsConcurrentExecutionAttempts {}该接口是一个空标记接口含义是该实现支持多个尝试同时执行见 SupportsConcurrentExecutionAttempts.java 源码注释。它适用于三类 SinkSinkSinkV2 / 新版 Sink APISinkFunctionOutputFormat。两个重要的边界规则任务级联生效如果任务中的任意一个算子不支持推测执行整个任务都会被标记为不支持推测执行。也就是说如果 Sink 不支持推测执行那么包含该 Sink 算子的任务将无法被推测执行。Committer 例外对于 Sink 实现Flink 会为 Committer 显式关闭推测执行——包括由 WithPreCommitTopology 和 WithPostCommitTopology 扩展出的算子。原因有二并发提交Concurrent Committing对不熟悉的用户可能引发意外问题而且 Committer 几乎不可能是批作业的瓶颈。如何验证推测执行的效果通过 Web UI 观察启用推测执行后当确实存在慢任务并触发了推测执行时在作业页面的顶点SubTasks标签页中可以看到推测执行尝试speculative attempts。在 Flink 集群的Overview与Task Managers页面上可以看到被屏蔽的 TaskManagerblocked taskmanagers。通过专用指标量化在 metrics.md 的 Speculative Execution 一节中定义了如下作业级指标仅在 JobManager 上可用Scope指标类型说明Job仅 JobManager 可用numSlowExecutionVerticesGauge当前时刻慢执行顶点的数量。Job仅 JobManager 可用numEffectiveSpeculativeExecutionsCounter有效的推测执行尝试数量即比其对应原始尝试更早完成的推测执行尝试数量。其中numEffectiveSpeculativeExecutions是衡量推测执行是否真正带来收益的关键指标推测尝试如果最终跑得比原始尝试还慢则属于无效推测不会计入该计数。建议在开启推测执行后结合这两个指标判断当前作业的慢节点情况与推测收益。小结推测执行通过慢任务检测 → 节点屏蔽 → 在健康节点上重部署新尝试 → 首个完成者胜出的机制缓解问题节点导致的批作业变慢。只需execution.batch.speculative.enabled: true即可开启但要求使用 Adaptive Batch Scheduler 与基于 DataStream API 的批作业。调优重点是两组参数调度侧的max-concurrent-executions/block-slow-node-duration检测侧的check-interval/baseline-ratio/baseline-multiplier/baseline-lower-bound。自定义 Source 若使用自定义 SourceEvent需让 SplitEnumerator 实现SupportsHandleExecutionAttemptSourceEvent自定义 Sink 需实现SupportsConcurrentExecutionAttempts才会参与推测执行。通过 Web UI 的 SubTasks 标签页、集群页面的被屏蔽 TaskManager以及numSlowExecutionVertices/numEffectiveSpeculativeExecutions两个指标可以直观评估推测执行的效果。【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表