ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

推测执行详解:从Hadoop MapReduce到Spark的调优实战

推测执行详解:从Hadoop MapReduce到Spark的调优实战 1. 一次真实的集群“掉队”事故我为什么开始重视推测执行大概两年前的这个时候我负责的一个离线数仓集群出了个诡异现象每晚跑核心ETL任务整个DAG都跑完了就卡在最后几个MapReduce job上。点开Hadoop Application页面一看其他Map任务早就完成就剩两三个任务卡在99%跑了一个多小时还没结束CPU使用率却低得可怜。当时第一反应是“数据倾斜”于是把热点key拆了、加桶、调并行度折腾半天下一次跑又换了个新的task卡住完全随机。后来才意识到我一直在解决错误的问题。这不是数据倾斜而是典型的推测执行缺失导致的“木桶效应”——集群里有几台老旧的物理机磁盘IO抖动严重某些task分到那些节点上就是跑不动而整个Job必须等最后一个task完成才结束。一个task拖慢几百个已经完成的task都在干等资源白白浪费SLA被一次次击穿。那次之后我把推测执行从头到尾捋了一遍包括它在MapReduce和Spark两代引擎里的实现差异、参数怎么调、哪些场景不能开、线上误判怎么排查。这篇文章就围绕这套实战经验来写。不管你是在维护自建Hadoop集群还是在用Spark做实时批处理搞懂推测执行的工作边界和调优逻辑处理这种“掉队任务”问题会更有底。2. 推测执行到底在解决什么问题一个关于“木桶短板”的分布式难题2.1 为什么分布式计算最怕“掉队任务”分布式计算的基本思路是“分而治之”把一个大数据集切分成多个分片分发到不同节点上并行处理。在没有故障的理想世界里每个task的耗时应该差不多整个Job的完成时间约等于大部分task的耗时加上调度开销。但真实集群不是理想世界节点异构、磁盘老化、网络抖动、内存争抢甚至同一台机器上别的业务在跑IO密集型的任务都会让某些task明显慢于其他task。这些拖后腿的task业内叫straggler也就是掉队任务。一个Job里有几百上千个task只要有一个掉队整个Job的完成时间就被它拉长了。更麻烦的是掉队任务往往是随机出现的你没办法提前预判哪个节点会跑得慢。指望运维把每台机器都修到性能一致既不现实也不经济。举一个很直观的例子一个MapReduce Job有1000个Map task每个task正常情况跑2分钟集群资源充足所有task并行执行Job大概3分钟完成。但如果99.9%的task都在2分钟跑完剩下1个task因为所在节点磁盘故障导致写map中间结果特别慢跑了30分钟才结束整个Job就被这个task拖到了30分钟。这时候加再多资源也快不了因为瓶颈在掉队任务本身。2.2 推测执行的核心思想让“替补队员”上场推测执行的设计思路并不复杂系统检测到某个task长时间没有完成并且进度明显落后于同批次的其他task就在另一台空闲节点上启动一个同样的task作为副本让这两个task同时跑。哪个先完成就采用哪个的结果后完成的直接杀掉回收资源。这有点类似足球比赛里的替补席——场上的主力队员状态不好教练不直接换人因为不确认是不是真的状态不好也可能是对手太强而是让替补开始热身准备。如果主力找回状态替补继续坐着如果主力确实不行替补立刻顶上。用比较小的资源代价换取整个Job完成时间的确定性。在MapReduce里这套机制是默认开启的对应的参数是mapreduce.map.speculative和mapreduce.reduce.speculative默认都是true。在Spark里推测执行同样默认开启参数是spark.speculation默认true。这里有个容易混淆的点Spark的推测执行机制和MapReduce不完全一样后面我会单独讲。2.3 一个关键前提推测执行只能在“多副本无副作用”的任务上使用推测执行能安全运行有一个隐含前提这个task必须是幂等的或者至少是“杀掉重跑不会产生副作用”的。Map task的处理结果是中国结果写到本地磁盘或内存由Reduce task拉取杀掉一个Map task的副本不会影响最终数据的正确性。Spark里的ShuffleMapTask同理输出的是给下游stage用的中间数据多个副本写同一份数据可能有点浪费但不会导致结果出错。而Spark的ResultTask直接产出最终结果的task就需要小心了如果它的计算逻辑里有写数据库、写外部文件这类副作用操作推测执行可能会触发两次写入。所以在生产环境里如果你的Spark作业涉及外部系统写入建议把含有副作用逻辑的stage单独处理或者考虑关闭推测执行避免重复写入带来脏数据。3. 从Hadoop到Spark两代引擎的推测执行机制差异3.1 Hadoop MapReduce按进度率判断“掉队”MapReduce的推测执行判断逻辑相对简单粗暴。ApplicationMaster会定期收到各个task的进度报告通过对比当前task的进度和同Job里其他task的平均进度计算出进度率。如果一个task的进度率明显低于平均值就判定为掉队启动推测执行。这里面有几个关键参数参数名默认值作用mapreduce.map.speculativetrue是否对Map task启用推测执行mapreduce.reduce.speculativetrue是否对Reduce task启用推测执行mapreduce.job.speculative.slowtaskthreshold1.0掉队task的进度率阈值低于平均值的这个倍数就触发mapreduce.job.speculative.slownode.threshold1.0掉队节点的阈值节点整体平均进度率低于这个倍数触发mapreduce.job.speculative.retry-after-no-speculate1000关闭推测执行后多久重新开始检测毫秒mapreduce.job.speculative.retry-after-speculate2000启动推测执行后多久再进行下一次检测毫秒举个例子帮你理解这些参数假设一个Job有10个Map task其中9个跑了50秒进度率大概是每秒2%剩下1个跑了50秒才完成20%进度率只有每秒0.4%远低于平均值的1倍阈值系统就会判定它为掉队task在另一台节点上启动一个相同task。新启动的task如果15秒跑完了系统直接杀掉原来的慢task采用新task的结果。这里有个比较反直觉的点推测执行并不是越快启动越好。如果启动太快可能一个task只是短暂波动副本启动后original task又恢复了这就白白浪费了一份计算资源。MapReduce里的retry-after-speculate参数默认2000毫秒意思是每次启动推测执行后至少要等2秒才做下一轮判断避免频繁误判。3.2 Spark基于“运行时间中位数”的推测逻辑Spark的推测执行实现和MapReduce差异挺大。Spark Driver里的TaskScheduler会定期检查所有正在运行的task如果发现有task的运行时间超过了同stage里已完成task运行时间的中位数的一定倍数并且剩余时间估算也超过阈值就触发推测执行。关键参数如下参数名默认值作用spark.speculationtrue是否启用推测执行spark.speculation.interval100ms检查频率Driver每隔多久扫描一次task状态spark.speculation.multiplier1.5运行时间超过已完成task中位数的多少倍触发推测spark.speculation.quantile0.75参考分位数默认取已完成task的75分位运行时间spark.speculation.quantile.duration未知实际按环境而定我来解释一下这几个参数的配合逻辑。quantile决定“参考基准”默认0.75表示取已完成task运行时间排序后的75分位数比如100个task完成了按运行时间从短到长排序第75个task的运行时间就是基准值。multiplier决定“超出比例”默认1.5表示当一个task的运行时间超过基准值的1.5倍时认为它是潜在的掉队task。这里要注意几个容易踩的坑quantile取0.75意味着必须有足够多的task完成才能算出中位数。如果stage的task数量太少比如只有几十个或者任务本身运行时间很短几秒钟就跑完推测执行可能还没来得及触发整个stage就已经完成了这时候推测执行完全不起作用。反过来如果集群资源非常紧张每个task都在排队等待资源运行时间普遍偏长推测执行可能会判定很多task都是掉队task然后启动大批副本导致资源竞争加剧反而拖慢整个Job。3.3 两种机制的直观对比我把两套机制放在一张表里对比方便理解它们的核心差异对比维度Hadoop MapReduceSpark判断基准同Job内所有task的平均进度率同Stage内已完成task的75分位运行时间 × 1.5检查频率由retry-after-speculate控制默认2秒由spark.speculation.interval控制默认100ms触发条件task进度率低于平均值一定倍数task运行时间超过基准值一定倍数副作用控制天然安全Map结果可覆盖ResultTask需谨慎可能重复执行副作用逻辑资源开销相对保守检查频率高资源开销相对更大核心区别在于MapReduce审视的是“速度”Spark审视的是“绝对耗时”。前者天然适应不同task之间执行时间的波动后者更适合“大家都差不多就你特别慢”的典型场景。实际使用中我一般根据Job的特点来选择要不要调整这些参数。4. 什么时候该开、什么时候该关推测执行的适用边界4.1 最适合开启推测执行的场景推测执行不是银弹它对场景有明确偏好。根据我的实际经验以下情况开启推测执行收益会比较大。场景一大量短task.map端作业。比如Hive里跑一个简单的count、filter、join前的预处理几千个Map task每个task处理几百MB数据耗时1-3分钟。这种场景下任何一台节点抖动都会让几个task明显变慢。而task之间相互独立推测执行成本很低——多跑几个副本撑死多消耗几个CPU核但换来的是整个Job不会因为某台物理机故障而延迟半小时。Spark里最常见的ETL清洗任务也属于这一类shuffle之前的Map阶段task数量多、单task耗时短推测执行收益非常明显。场景二节点异构明显的集群。如果你和我一样集群是分批采购的有新的物理机也有用了五六年的老机器那么异构节点的性能差异会直接体现在task执行时间上。老机器的CPU主频低、磁盘IO慢同样的数据量跑起来就是比新机器慢一半。这种时候靠推测执行做“劣后淘汰”比强制让所有task都在同构节点上跑要划算得多。场景三SLA要求严格、任务失败容忍度低的场景。比如每天凌晨必须完成的数据对账任务或者业务方强依赖的T1报表任务宁可多消耗10%的资源也要保证不会因为个别节点抖动导致整体延迟。这种场景下推测执行可以看作是为“失败确定性”买的保险。4.2 强烈建议关闭推测执行的场景场景一每个task都会写外部系统。比如SparkStreaming写HBase、写MySQL、写Kafka或者Spark批处理里最终stage直接写业务库。如果启用了推测执行同一个task的两个副本都可能执行写操作极端情况下会造成主键冲突、重复写入。虽然有些系统做了幂等比如HBase的put是覆盖写但大多数业务系统做不到所以这种作业我会直接关闭推测执行。场景二数据倾斜本身就很严重的Job。很多人没意识到“推测执行 数据倾斜”会叠加出灾难。数据倾斜意味着某些task天然比其他task多处理几倍的数据运行时间天生就长。推测执行会把这些“正常地慢”的task误判成掉队task启动大量副本。这些副本同样要处理倾斜数据依然跑得慢最终结果是集群资源被副本占满真正健康的task反而排队等待资源整个Job比不开推测执行还要慢。所以我处理倾斜问题时第一步永远是关推测执行先解决倾斜再决定要不要恢复开启。场景三task执行时间普遍极短秒级的作业。比如一些轻量级的SparkSQL查询task几十秒就完成Driver还没来得及做第二次speculation.interval检查stage已经跑完了。这时候推测执行不仅没收益反而因为Driver频繁检查task状态、计算中位数给Driver带来额外压力。遇到这种作业直接关掉反而干净。场景四集群资源非常紧张没有空闲slot。推测执行本质上是用“冗余资源”换“时间确定性”。如果集群已经满负载运行每启动一个推测副本意味着某个正常task要排队等资源。这种内耗可能导致整个Job的吞吐量反而下降。我的经验是集群负载持续超过70%的时候推测执行的收益已经大打折扣优先考虑扩容而不是指望推测执行解决问题。4.3 我的选型判断清单每次给一个新的Job配置推测执行我都会先问自己几个问题你以后也可以按这个思路来判断这个Job的task之间有没有共享状态如果有关掉。有没有外部系统写入如果有关掉或者确保幂等。task执行时间分布是否均匀如果天然就不均匀比如直方图分布很宽警惕“误杀”。集群还有没有空闲资源没有空闲资源就别开效果为零甚至负。这个Job的SLA敏感度有多高高SLA任务可以承担少量资源浪费低SLA任务优先节约资源。5. 推测执行“误杀”时刻掉队、倾斜和快速任务的三方博弈5.1 一个完整的故障排查链路从“莫名开启的副本”说起有一次线上Spark任务突然变慢我一看Spark UI发现某个stage启动了比往常多一倍的task数量。点进详情页看到大量task被标记为Speculative也就是说推测执行被频繁触发。当时的任务是一个简单的ETL清洗没有任何外部系统写入理论上推测执行开启也没问题但结果却变慢了。我的排查思路是这样的第一步先看task执行时间分布。发现有个别task的运行时间确实是其他task的2-3倍但这部分task处理的数据量也是其他task的2-3倍。也就是说掉队的原因不是节点性能问题而是数据量差异——典型的数据倾斜。第二步确认倾斜的根本原因。通过Spark UI看每个task的输入数据量发现某个分区下有一个超大key按某个业务字段做join时所有相同key的数据都集中到一个分区里。这时候推测执行把倾斜task当成掉队task启动了多个副本每个副本都要处理同一份倾斜数据跑得一样慢。副本之间抢资源正常task反而被挤到后面。第三步处理方式先关闭推测执行避免资源进一步浪费然后用加盐salting的方式把大key拆散重新跑性能立刻恢复。这次事故给我留下很深的印象推测执行的触发逻辑“看起来合理”但它不会判断task“为什么慢”只会判断“你是不是慢了”。如果你没有把自己的数据特征搞清楚推测执行就是一把乱挥的刀。5.2 快速任务场景下的“误杀”比想象中更容易发生还有一个高频误杀场景是task运行时间极短的作业。比如你有一个Spark Streaming任务micro-batch处理时间只有几秒。Spark的speculation.interval默认是100ms理论上Driver每100ms就会检查一次task状态但实际判断时还需要quantile也就是需要一定数量的已完成task来算中位数。如果batch处理时间短stage瞬息万变中位数还没算出来任务已经结束了。另一种更隐蔽的误杀发生在task数量很少的时候。假设一个stage只有10个task其中1个task因为节点网络抖动慢了1倍其他9个都是正常速度。按Spark的逻辑已完成task的75分位时间作为基准正常的9个task已经完成了算出中位数然后那个慢task开始被怀疑。问题在于5个乃至4个task的情况下中位数受异常值影响很大可能把正常偏慢的task也误判为掉队。这也是为什么我不太建议在task数量极少的作业里开推测执行的原因。5.3 一个被忽略的操作结合节点黑名单判断在实际排查掉队task时我总结出一个规律掉队task往往集中在少数几台节点上。这是判断到底是“节点故障”还是“数据倾斜”的一个很好的辅助手段。如果多个Job里掉队的task都落在同一批节点上比如那几台老机器、磁盘告警的机器那基本可以确定是节点问题推测执行确实能解决。这时候更优的做法是把这些节点加入黑名单yarn.resourcemanager.nodes.exclude-path让调度器不再往这些节点上分配task从源头上解决。如果掉队task分布的节点很随机但每个task处理的数据量差异很大那优先怀疑数据倾斜推测执行解决不了问题甚至帮倒忙。这两类情况需要分开处理不能一上来就调推测执行参数否则只是治标不治本。6. 线上调优实操这几个参数我踩过坑建议你这样调6.1 Hadoop MapReduce侧参数配置在mapred-site.xml里我实际用下来比较稳妥的配置组合是这样的property namemapreduce.map.speculative/name valuetrue/value /property property namemapreduce.reduce.speculative/name valuetrue/value /property property namemapreduce.job.speculative.slowtaskthreshold/name value1.2/value /property property namemapreduce.job.speculative.slownode.threshold/name value1.2/value /property property namemapreduce.job.speculative.retry-after-no-speculate/name value2000/value /property property namemapreduce.job.speculative.retry-after-speculate/name value4000/value /property我把slowtaskthreshold从默认的1.0调到了1.2意思是task必须有1.2倍的差距才会判定为掉队。这是为了给正常的task波动留出余量减少误杀。把retry-after-speculate从2000毫秒调到4000毫秒降低检查频率避免在task波动时频繁启动副本。这里有个参数组合的细节slowtaskthreshold和slownode.threshold是同时生效的一个task只有同时满足“task进度率低于平均值1.2倍”和“所在节点平均进度率低于集群平均值1.2倍”才会触发推测。所以如果只是想放宽单个task的判断条件只调slowtaskthreshold就够了slownode.threshold不要随便调否则会影响node-level的判断逻辑。6.2 Spark侧参数配置在spark-submit脚本里我常用的推测执行配置是这样spark-submit \ --conf spark.speculationtrue \ --conf spark.speculation.interval1000 \ --conf spark.speculation.multiplier1.8 \ --conf spark.speculation.quantile0.8 \ --conf spark.speculation.quantile.duration15000 \ ...解释一下为什么这么调spark.speculation.interval默认是100ms我建议调到500ms或1000ms。原因很实际如果集群规模比较大Driver本身就要处理很多task状态信息每100ms扫描一次所有正在运行的taskDriver的压力会明显增加。调整到1000ms后对检测灵敏度的影响其实很小因为大多数掉队task都是分钟级问题100ms的检查精度没什么必要。spark.speculation.quantile默认是0.75multiplier默认是1.5。我习惯把multiplier调高到1.8quantile调到0.8。这样触发条件变严格只有明显掉队的task才会被判定。尤其当集群负载本身偏高时太敏感的触发条件会放大资源竞争问题。spark.speculation.quantile.duration是另一个重要的参数它表示task运行时间达到多少毫秒才参与中位数计算。默认值是0意思是所有task都被纳入计算。如果stage里有大量秒级短task它们的运行时间会拉低中位数导致正常范围内的task也容易被判定为掉队。我一般设置为15000即运行时间超过15秒的task才作为参考样本过滤掉短task的干扰。还有一个容易被忽略的参数spark.speculation.minTaskRuntime部分版本里存在用于设定task的最小运行时间低于这个值的task不会触发推测。这个参数和quantile.duration功能有重叠如果你的版本里没有quantile.duration可以用这个参数替代。6.3 动态调整不同Job用不同配置的实践实际生产里我不建议全局统一配置推测执行参数。同一个集群上跑的作业五花八门有的适合开有的需要关。最理想的方式是在spark-submit时按作业类型传入不同配置让推测执行的开关和灵敏度跟着作业特征走。我的做法是维护一个“作业配置清单”里面按作业类型记录推荐的推测执行设置。比如作业类型推测执行推荐配置常规ETL清洗无外部写入开启multiplier1.8, quantile0.8, interval1000数据倾斜常见报表临时关闭先解决倾斜问题再评估开启写HBase/MySQL/外部接口关闭避免重复写入实时流处理关闭短task场景收益低增加Driver压力长时跑批任务30分钟以上开启可按默认配置适当加大interval如果你是在维护一个共享Hadoop集群可能无法为每个作业定制参数那就至少做到在mapred-site.xmlorspark-defaults.conf里设置一个相对保守的默认值避免误杀大面积发生。7. 从推测执行到自适应调度我的一些延伸思考7.1 一个和推测执行很像但不是一回事的概念动态分区有读者可能把推测执行和动态分区搞混其实它们是不同层面的机制。推测执行解决的是“task掉队”问题动态分区解决的是“分区内数据不均匀”问题。前者是在task调度层面做文章后者是在数据切分层面做文章。两者可以同时使用但在数据倾斜场景下动态分区通常更治本。7.2 新引擎里推测执行的演化和替代方案如果你在关注比较新的计算引擎会发现推测执行在下一代系统里有升级趋势。比如Spark 3.x开始支持基于DAG的Adaptive Query ExecutionAQE它会根据shuffle后的实际数据分布动态调整分区数从而缓解一部分数据倾斜问题进而减少掉队task的产生。虽然这不能完全替代推测执行但它确实在“从源头减少掉队发生的概率”。另外一些云厂商的Serverless Spark服务会默认关闭推测执行理由是Serverless环境通常对资源隔离做得更好节点性能更均匀掉队概率低而推测执行会带来额外成本。这个选择本身也验证了一个观点推测执行是“在不可靠基础设施之上做的防御机制”基础设施越可靠推测执行的价值越低。7.3 我对推测执行的最终态度做了这几年分布式计算我越来越觉得推测执行像是一个“必要的妥协”。它不可能完美区分“节点慢”和“数据歪”因为对于一个分布式系统来说它能够利用的信息永远是有限的。它存在的意义不是消灭掉队而是在不确定性的环境下用可接受的资源浪费换取整体进度的确定性。所以我的建议是别再纠结默认参数是不是最合理的先花点时间了解你的集群和你的数据特征。真正稳定运行的作业往往不是靠一个完美参数调出来的而是靠理解每个机制背后的假设条件然后在合适的场景里做合适的选择。如果你现在正被某个“卡在99%”的任务坑得焦头烂额我建议你先按这套思路操作一遍打开Spark UI或者Hadoop Application页面看看掉队task的分布和输入数据量判断到底是节点问题还是数据问题然后再决定推测执行的开关和参数。这样做一次你对推测执行的理解会比看十篇文档都深刻。
返回列表