ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

分布式作业为什么越跑越慢?推测执行机制与Spark调优实战

分布式作业为什么越跑越慢?推测执行机制与Spark调优实战 做分布式计算的人多少都被这种场景折磨过一个几千并发的Spark任务跑得好好的进度条稳步往前走突然就到了99%卡住不动。其他所有task早就跑完了就剩那么一两个还挂在Executor上整个Stage跟着一起拖最后跑出翻倍甚至好几倍的总耗时。有人第一反应是数据倾斜上去看数据分布也有人脱口而出“把推测执行打开试试”。这个“推测执行”到底是个什么机制为什么能跟数据倾斜并列成为分布式作业性能排查的第一批怀疑对象网上讲原理的不少但真正从配置到实战踩坑讲清楚的比较少。这篇文章我用实际生产环境的经验把推测执行从头到尾拆一遍从它要解决的问题、内部判断逻辑、参数调优到什么时候该关、什么时候该开一次说透。适合正在做Spark、MapReduce调优或者写分布式任务被“慢任务”困扰的工程师参考。1. 为什么慢任务能拖垮整个分布式作业1.1 慢任务从哪来分布式计算框架的基本调度单位是task一个Spark Stage会被拆成若干个互不相干的task分发到不同节点上并行跑。理想状态下每个task拿到的数据量差不多、跑的时间差不多整个Stage的总耗时约等于单个task的平均耗时。但现实里这个理想状态几乎不存在因为拖后腿的task来源太多了。首先是硬件差异。云厂商的机器看着规格一样实际底层CPU型号、磁盘类型、网络带宽可能完全不同。共享型实例更夸张邻居虚拟机一吵你的CPU时间片直接缩水30%。同样的task在好的机器上跑了1分钟在差的机器上可能跑了4分钟。其次是资源争抢。同一台物理机上跑了多个Executor或者同一批机器上还有其他团队的任务CPU、磁盘IO、网络带宽都被分摊。表现就是任务偶尔出现很长的GC停顿或者shuffle阶段的数据拉取极其缓慢。然后是最经典的数据倾斜。某些task分配到的分区数据量远大于其他分区比如用户维表关联的时候某个超大key占了大头这种情况下不是机器慢是task本身的工作量就重。注意区分机器慢是task本身不重但运行慢数据倾斜是task本身即使跑得动也要花那么久。还有一个大家容易忽略的某个task运行所在的节点出了问题比如磁盘坏道反复重试、网络闪断频繁重连。这种task的进度可能一直卡在某个百分比不动但框架还没判定它失败其他task早就结束整个作业就卡在这一个“半死不活”的任务上。1.2 等一个“掉队者”的代价很多人觉得一个task慢而已最坏情况就是总耗时被拉长一点能拉长多少这个账得仔细算。MapReduce模型里一个reduce task要等上游所有map task的输出都拉齐了才能开始merge和reduce。Spark虽然每个Stage内部task相互独立但Stage之间有宽依赖下游Stage必须等上游Stage的全部task做完才能启动。也就是说只要有一个task不结束整个Stage就不会结束下游Stage即使有资源空闲也得干等。这个“等”的代价在集群规模大的时候会被指数级放大。假设一个Stage有1000个task平均每个task跑2分钟正常情况下2分多钟结束。但如果有1个task因为机器问题跑了20分钟整个Stage的完成时间就是20多分钟10倍差距就这么出来了。更麻烦的是大集群里这种异常task几乎必然出现——节点越多某个节点出幺蛾子的概率就越高。这也是为什么早期单一机器跑MapReduce没这个烦恼上了大规模集群之后才必须处理掉队者问题。有人会说那直接把超时的task杀掉重跑不就行了问题在于框架怎么判断“这个task是真的有问题”还是“本来就该跑这么慢”。如果贸然杀掉一个只是数据量偏大的task重跑一遍照样要花同样甚至更长的时间还得搭上调度重试的额外开销。所以需要一种机制——不是简单杀掉重跑而是在不中断原任务的前提下“再找一个副本并跑”谁先跑完谁说了算。这就是推测执行的核心思路。2. 推测执行的工作原理2.1 机制拆解“备胎任务”如何启动推测执行的完整英文是Speculative Execution国内也叫“投机执行”。我更喜欢叫它“备胎机制”——正选任务跑着如果框架觉得它不对劲就再找个空闲节点起一个一模一样的副本任务。两个任务跑同一个分区的数据谁先成功提交谁的结果就是最终结果另一个任务直接被框架kill掉。执行过程拆开看大致分成四步第一持续监控。框架会定期收集所有task的运行时指标包括启动时间、运行时长、当前进度、已处理数据量等。这几个指标是后面所有判断的依据。第二识别可疑任务。这就是整个机制的核心判断逻辑我会在2.2里细说。简单来讲框架会给每个task算一个“预期完成时间”如果某个task按照当前运行速度预计完成的时间比同一个Stage里其他同类task的平均完成时间晚很多就会被打上“可疑”标签。第三启动副本。可疑任务不会马上被替代而是等待进一步确认。确认之后框架会在有空闲资源的节点上启动一个完全相同的task运行同样的分区数据。第四结果竞争。正选和副本两个任务同时跑谁先成功完成框架就以它的输出为准。后完成的任务即使最终也跑成功了结果也会被丢弃。这里有一个很多人没意识到的好处推测执行不是把慢任务停了再重启而是让两个任务并行跑其中快的那个能拿到最终结果。哪怕最后正选任务赢了副本也已经消耗了资源但这个作业的完成时间被“救”回来了——这笔账在作业整体被拖累的代价面前通常是划算的。2.2 什么时候才判定“需要推测”判断逻辑这个细节值得细讲因为它直接决定推测执行是帮手还是祸害。如果判断得太激进满集群都是重复跑的任务资源被白白吃掉判断得太保守又起不到救场效果。Spark的实现参考了Hadoop MapReduce的思路核心判断量有两个一个是推测执行开始后被监控任务的“平均进度”和“已运行时间”另一个是同一个Stage内所有成功完成任务的“平均进度”和“平均运行时间”。具体判定时用一个阈值倍数做比较——Spark的spark.speculation.multiplier默认是1.5spark.speculation.quantile默认是0.75。意思是当集群里已完成task数量达到Stage总量的75%某个正在运行的task的预计完成时间如果超过已完成task平均完成时间的1.5倍这个task就会被标记为可疑候选进而被启动副本。为什么强调“已完成任务达到75%”我的理解是需要足够多的已完成任务样本才能算出合理平均值。如果Stage刚开始跑统计样本太少一个task稍微慢一点就被误判推测执行就会频繁误触发。这个设计其实挺巧妙先跑完足够多任务有了“大多数人都这个速度”的基准再来揪出真正掉队的那个人。MapReduce的判定则更侧重ExecutionTimes的统计分布和平均完成比例。YARN的TaskScheduler会对每个task记录运行时间和Progress然后周期性地检测是否存在运行时间远超过其他任务、且进度明显落后的task。判断过程中还会排除一些特殊情况比如任务启动时间太短还没热身不会上来就误判。整体来看推测执行的共同逻辑可以归纳成一句话在大多数同类任务都已经完成后衡量剩下还在跑的任务是否“显著落后”。这个“显著落后”的标准其实就是经验值这也是为什么推测执行参数跑不同作业会有完全不同的最优值。2.3 副本任务跑重复了怎么办“两个task同时跑同一个分区的数据最后提交结果的时候怎么保证不会乱套”这是很多人第一次听说推测执行会问的问题。实际上不用太担心一致性因为框架层面对task的状态管理是有强力约束的。Task在有向无环图DAG里的角色是生产者。对于仅影响内部状态的task比如某个Stage内部的map类task它的输出是中间数据谁先完成谁就把数据写回给Driver或ShuffleService后完成的那个task即使也写了一份框架会通过任务ID的归属判断把多余的结果丢弃。对于影响最终结果的taskexecute一次只能有一个提交状态转换Executor会先向Driver汇报成功Driver确认任务状态从RUNNING变为SUCCESS之后其他迟到的完成通知就不会再改变状态了。拿Spark的Shuffle来说每个map task会生成数据文件和索引文件文件的命名包含task的ID。如果两个task跑同一个数据分区最终只有获胜的那个task的ID对应的文件会被下游Stage读取。失败副本产生的中间文件要么被覆盖要么在stage清理时被删掉不会影响后续计算正确性。真正要操心的不是“会不会算错”而是“副本任务本身能不能跑起来”。推测执行要求集群里还有空闲资源来运行副本。如果集群已经是满载状态Driver给副本分配不到Executor推测执行就算触发了也是空转。这个问题在参数调优时很重要后面第3章会专门聊。2.4 MapReduce与Spark的实现差异两个主流引擎虽然思路一致实现取舍差别还挺大。Hadoop MapReduce的推测执行是默认开启的map和reduce阶段都可独立配置。老版本里有个被吐槽很久的问题reduce阶段推测执行的效果没那么好因为reduce task要从各个节点拉取map输出副本任务的网络IO开销可能比正选任务还大甚至引起重复拉取风暴后来很多团队在实践里干脆把reduce推测关掉。Spark侧则是默认关闭。spark.speculation默认值是false真要说为什么一个主要原因是Spark阶段内部task失败重试代价相对低另一个原因是Spark流式计算场景下推测执行会引入重复计算和结果延迟的不确定性宁可不开。但在批处理场景中尤其跑SQL分析、ETL这种任务开推测执行的收益通常明显大于成本。还有个差异在动态调整上。MapReduce几乎只支持通过配置文件开关Spark则提供了多个细粒度参数控制推测的条件和频率。另外Spark 3.0以后引入了Dynamic Partition Pruning和Adaptive Query Execution和推测执行搭配使用效果更好某种意义上AQE也抢了一部分推测执行该干的活——它通过聚合shuffle统计信息动态调整分区数直接缓解了“分区数据不均导致个别task过重”的问题。3. 参数配置与调优思路3.1 Hadoop MapReduce相关参数如果你的作业跑在Hadoop MapReduce上核心参数是下面这几个参数名默认值说明mapreduce.map.speculativetrue是否对map任务启用推测执行mapreduce.reduce.speculativetrue是否对reduce任务启用推测执行mapreduce.job.speculative.slowtaskthreshold0.3慢任务判定阈值根据执行时间均值与方差计算mapreduce.job.speculative.slownodethreshold1.0慢节点判定阈值用于排除节点整体性能差的情况mapreduce.job.speculative.speculativecap0.1推测执行任务数的上限比例防止过度推测其中speculativecap值得拿出来说。它在框架层面对同时启动的推测副本数做了上限控制防止一个阶段里一大半任务都在跑副本资源被疯狂占用。运维经验是老集群如果整体资源紧张建议把reduce的推测执行关掉只留map的。因为reduce阶段一旦开启推测重复的fetch、merge操作经常得不偿失。3.2 Spark相关参数Spark的推测执行参数集中在spark.*前缀下最常用的就几个参数名默认值说明spark.speculationfalse总开关spark.speculation.interval100ms扫描间隔每隔多久检查一次有没有可疑taskspark.speculation.quantile0.75已完成task比例达到多少才开始判定spark.speculation.multiplier1.5慢任务倍率阈值spark.speculation.quantile的结合逻辑见上文判断算法核心第一次看spark.speculation.interval默认100ms的时候我挺意外这个间隔实在太短了。过短的扫描间隔意味着Driver要频繁地做检测计算任务一多反而给Driver增加负担。我实际用的环境中一般调到1000到2000ms这样既不会漏检也不至于反复做无意义的扫描。特别是Stage里有几千个task时每次检测都要遍历所有正在运行的task并做进度预估频率太高压力不小。另外还有一个容易踩坑的地方当你用了Spark的动态资源分配spark.dynamicAllocation.enabledtrue并且设置spark.dynamicAllocation.executorIdleTimeout较短时如果任务长时间“没东西干”等待被释放一旦推测执行被触发要启动副本却发现Executor额度已经被缩减了副本启动就要排队。建议是在批量作业中开推测执行时先检查一下动态分配的配置是否有冲突。3.3 什么场景该开什么场景不该开直接上结论我现在的经验判断规则是这样的建议开启的场景作业运行时长较长Task数量大集群规模在几十台以上运行环境硬件异构明显比如混用物理机、虚拟机和抢占式实例上游数据随机性大经常出现个别分区内容特别复杂Stage主要是CPU密集类型数据倾斜不明显慢主要来自节点性能差异建议关闭的场景作业本身很轻几十秒就跑完了。推测执行还没完成判断作业早就结束了开了纯属多余数据倾斜是主要矛盾。倾斜导致的慢task你用推测执行给副作用节点加再多的副本数据量不均这个问题本身未被解决反而白跑集群资源利用率常年跑在90%以上没有空余资源跑副本。推测执行触发后副本拿不到资源同时正选任务还被抢占偏实时或流式的场景一次性事件的重复计算成本不允许还有一个特别重要的使用习惯推测执行不是SQL里的索引加上就能让慢查询变快。它是一个兜底机制作用是防止“个别task的异常拖垮整个作业”而不是主动优化task的执行效率。如果你发现一个作业频繁触发推测说明作业本身存在系统性的慢任务问题应该去定位根因——是数据倾斜还是Executor异常的节点过多——而不是单纯依赖推测执行续命。4. 实操配置案例与常见坑4.1 一个Spark作业开启推测执行的配置示例直接给一份我在生产环境常用的Spark Submit配置片段仅供参考版本基于Spark 3.xspark-submit \ --master yarn \ --deploy-mode cluster \ --conf spark.speculationtrue \ --conf spark.speculation.interval1000ms \ --conf spark.speculation.quantile0.75 \ --conf spark.speculation.multiplier1.5 \ --conf spark.sql.adaptive.enabledtrue \ --conf spark.sql.adaptive.coalescePartitions.enabledtrue \ --class com.example.etl.DailyETL \ daily-etl.jar这套配置的思路是打开总开关扫描间隔放宽到1秒保持默认的0.75和1.5阈值同时打开AQE做动态分区合并。对于大部分批处理作业来说这是一个首版配置跑一圈看看任务日志里speculation的输出和事件时间线再往下调。需要注意spark.speculation只对通过SparkContext提交的job有效。如果你用Spark SQL的thrift server跑查询这个参数要在服务端配置里设置不能在SQL会话里一条命令搞定。4.2 推测执行和Spark Web UI的配合观察开了推测执行之后想在UI上确认有没有生效关键看两个地方。第一个是Executor的Task列表中是否有Task ID相同的两行记录。这两行的“Launch Time”不同一个是最初启动一个是副本启动。正常情况下两个task最终一个标记为SUCCESS一个标记为KILLED。第二个是Stage详情页的“Skipped Tasks”或“Dead Tasks”统计。虽然KILLED不直接等于推测执行但在开启speculation的作业里如果发现大量KILLED任务且杀掉的都是同一批分区ID的任务那基本就是推测执行在起作用。这里有个非常容易误读的细节KILLED任务并不意味着作业失败而是说明有一个“备胎”赢了另一个被淘汰。不要一看到KILLED就以为出故障了。反过来如果KILLED任务数量极其密集地集中在一个Executor节点上那说明这个节点的机器大概率有问题推测执行是在反复救火应该考虑把这个节点拉黑或去医院检查硬件。4.3 推测执行对数据倾斜误判的例子我来说一个踩过的坑。曾经有一个ETL任务某个表关联后做聚合个别用户的数据量特别大分配到这个用户所在分区的task要跑40分钟其他task只需要5分钟。开推测执行之后80%的task完成时那个跑40分钟的task被判定为可疑框架给它起了两个副本结果三个task同时抢一个分区的数据每个都跑了40多分钟。从作业总耗时来看原来最慢40分钟开了推测执行反而多消耗了大量计算资源总时间还更高。因为副本数量的增加让那台机器上出现了严重的IO争抢跑得更慢了。后来怎么改的并没有靠推测执行而是针对那个大用户做了单独的key加盐打散然后在聚合阶段再合并回来。数据倾斜类的慢task正确解法是调整数据分布让每个分区的数据量尽可能接近而不是起一堆副本去抢同一块已经很大的“蛋糕”。从那以后我给自己定了个原则日志里出现某个partition的任务连续多次被推测执行那就别再看参数了先拉Hive或者数据源的数据分布图绝大多数都是倾斜问题。4.4 资源水位与推测执行的关系推测执行是一个“额外消耗资源换时间”的手段你每一份资源都去跑副本正选任务和副本的任务其实在做同样的计算。这意味着集群的空闲资源越充足推测执行效果越好。我有一次在资源利用率大概95%的集群上给一个跑批作业开了推测执行结果任务反而变慢了。原因很简单副本任务起来之后和正选任务抢同一个CPU核两个都变慢了最后“谁先跑完”变成了“看谁都没跑完”。推测执行适合集群有缓冲余量的情况有些公司甚至专门为推测执行保留了5%~10%的“缓冲资源池”。如果你的作业队列已经非常拥挤建议评估是不是先扩容或者先等队列空闲了再跑而不是指望推测执行在满载集群上创造奇迹。4.5 并行度过低时的特殊处理有时候任务慢不是因为机器性能而是Stage的并行度本身设置得太低。比如一个Stage只有几个task每个task处理一个巨大的分区。这个时候spark.speculation.quantile0.75会导致判断迟迟不触发——因为总共就三四个task完成75%也就是至少完成2个才开始测算如果剩下一个跑太久前面两个task的完成时间平均值反而被拉低。这种场景下单纯开推测执行效果通常不好。正确的做法是先看输入数据量调整spark.sql.shuffle.partitions或者spark.default.parallelism让并行度符合资源的实际承载能力。推测执行应该作为并行度优化之外的补充手段而不是替代品。5. 常见问题与排查技巧实录5.1 典型问题速查表现象可能原因处理建议开启推测后总耗时没有减少反而上升集群资源利用率过高副本抢占正选资源关闭推测降低提交作业并发数日志中大量task被KILLED推测执行机制正常发挥作用检查KILLED是否集中在同一分区确认是否存在数据倾斜推测执行始终不触发interval过长或quantile比例过高调短扫描间隔或降低quantile到0.5~0.75KILLED的任务都集中在同一台机器该机器硬件异常排查节点CPU/内存/磁盘考虑中从资源池移除Spark任务出现重复数据结果大概率不是推测执行造成检查Shuffle重试机制和Task重试配置推测执行不会导致结果重复开启后Driver GC压力变大interval过短导致频繁扫描调大interval到1000ms~3000ms5.2 如何区分慢任务是推测执行还是数据倾斜判断依据不复杂看日志里慢任务的input数据量就够了。如果Web UI上显示的Shuffle Read或Input数据量和别的task差不多但运行时间高出数倍那么推测执行通常可以解决。如果数据量本来就比其他task大好几倍这就是倾斜去优化数据分布才是正道。再补充一个更实战的判断法当Stage跑到60%~70%左右时进去看一眼此时慢task和已完成task的进度差。如果进度差的趋势是“整体匀速、只是速度慢”等待和推测执行都行如果发现进度差是在某个百分比附近长时间不动比如卡在87%很久了那大概率是shuffle拉取或某个外部依赖出现问题推测执行未必能救。5.3 其他替代手段重试、黑名单与动态调整最后我想说推测执行不是解决慢任务的唯一手段甚至不总是最优手段。几个常见替代方案值得大家尝试。工作节点黑名单机制。YARN在任务反复失败后会将该节点标记为“黑名单”后续任务不再调度到它上面。如果你的集群问题节点比较稳定利用黑名单机制比每次靠推测执行救场更省资源。把supervisor的healthCheck做扎实。部分慢任务的根因是节点负载过高、磁盘IO固态衰退如果运维能提前把异常节点体检出来推测执行的压力自然小很多。自适应执行AQE。Spark 3.0以后带上了这个利器特能解决“分区数设置不合理”引发的个别task偏大问题。AQE会在shuffle之后动态合并小分区均匀化下游task的数据量建议和推测执行一起开两者定位不同一个管分布一个管兜底。动态资源分配协调。前面稍微提过开启动态分配时需要防止Executor被释放导致副本任务排队。实际操作中可以把spark.dynamicAllocation.maxExecutors适当调大给推测副本预留出空间。写在最后的一点体会推测执行这个机制用一句话总结的话就是在Face真实世界中“总有一个进程会掉队”的背景下用一部分冗余计算换作业的稳定性。分布式系统里节点故障和你不能完全杜绝网络抖动能让大多数作业不因为一两个异常task而“卡死”这个设计本身就非常实用。但它不是银弹甚至用错了还会添乱。我个人这几年运维下来最深的体会是不要一遇到慢任务就开推测执行。先开后看数据、先确认根因再动手推测执行才是一把好用的枪。如果你正准备在一个长期运行的批处理作业上尝试它建议先在一个Staging环境小规模跑一遍对比开和不开的总耗时与资源用量再用数据说话。最后分享一个小技巧如果你们集群有完善的任务监控可以给“同一Stage内task完成时间标准差”单独做一个告警指标。当标准差突然放大往往意味着慢任务开始冒头了这时候再决定打开推测执行或调整参数效果远比一直开着要好。祝大家跑批顺利。
返回列表