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.speculative | true | 是否对Map task启用推测执行 |
mapreduce.reduce.speculative | true | 是否对Reduce task启用推测执行 |
mapreduce.job.speculative.slowtaskthreshold | 1.0 | 掉队task的进度率阈值,低于平均值的这个倍数就触发 |
mapreduce.job.speculative.slownode.threshold | 1.0 | 掉队节点的阈值,节点整体平均进度率低于这个倍数触发 |
mapreduce.job.speculative.retry-after-no-speculate | 1000 | 关闭推测执行后,多久重新开始检测(毫秒) |
mapreduce.job.speculative.retry-after-speculate | 2000 | 启动推测执行后,多久再进行下一次检测(毫秒) |
举个例子帮你理解这些参数:假设一个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.speculation | true | 是否启用推测执行 |
spark.speculation.interval | 100ms | 检查频率,Driver每隔多久扫描一次task状态 |
spark.speculation.multiplier | 1.5 | 运行时间超过已完成task中位数的多少倍触发推测 |
spark.speculation.quantile | 0.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 MapReduce | Spark |
|---|---|---|
| 判断基准 | 同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要求严格、任务失败容忍度低的场景。
比如每天凌晨必须完成的数据对账任务,或者业务方强依赖的T+1报表任务,宁可多消耗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> <name>mapreduce.map.speculative</name> <value>true</value> </property> <property> <name>mapreduce.reduce.speculative</name> <value>true</value> </property> <property> <name>mapreduce.job.speculative.slowtaskthreshold</name> <value>1.2</value> </property> <property> <name>mapreduce.job.speculative.slownode.threshold</name> <value>1.2</value> </property> <property> <name>mapreduce.job.speculative.retry-after-no-speculate</name> <value>2000</value> </property> <property> <name>mapreduce.job.speculative.retry-after-speculate</name> <value>4000</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.speculation=true \ --conf spark.speculation.interval=1000 \ --conf spark.speculation.multiplier=1.8 \ --conf spark.speculation.quantile=0.8 \ --conf spark.speculation.quantile.duration=15000 \ ...解释一下为什么这么调:
spark.speculation.interval默认是100ms,我建议调到500ms或1000ms。原因很实际:如果集群规模比较大,Driver本身就要处理很多task状态信息,每100ms扫描一次所有正在运行的task,Driver的压力会明显增加。调整到1000ms后,对检测灵敏度的影响其实很小,因为大多数掉队task都是分钟级问题,100ms的检查精度没什么必要。
spark.speculation.quantile默认是0.75,multiplier默认是1.5。我习惯把multiplier调高到1.8,quantile调到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清洗(无外部写入) | 开启 | multiplier=1.8, quantile=0.8, interval=1000 |
| 数据倾斜常见报表 | 临时关闭 | 先解决倾斜问题,再评估开启 |
| 写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 Execution(AQE),它会根据shuffle后的实际数据分布动态调整分区数,从而缓解一部分数据倾斜问题,进而减少掉队task的产生。虽然这不能完全替代推测执行,但它确实在“从源头减少掉队发生的概率”。
另外,一些云厂商的Serverless Spark服务会默认关闭推测执行,理由是Serverless环境通常对资源隔离做得更好,节点性能更均匀,掉队概率低,而推测执行会带来额外成本。这个选择本身也验证了一个观点:推测执行是“在不可靠基础设施之上做的防御机制”,基础设施越可靠,推测执行的价值越低。
7.3 我对推测执行的最终态度
做了这几年分布式计算,我越来越觉得推测执行像是一个“必要的妥协”。它不可能完美区分“节点慢”和“数据歪”,因为对于一个分布式系统来说,它能够利用的信息永远是有限的。它存在的意义,不是消灭掉队,而是在不确定性的环境下,用可接受的资源浪费换取整体进度的确定性。
所以我的建议是:别再纠结默认参数是不是最合理的,先花点时间了解你的集群和你的数据特征。真正稳定运行的作业,往往不是靠一个完美参数调出来的,而是靠理解每个机制背后的假设条件,然后在合适的场景里做合适的选择。
如果你现在正被某个“卡在99%”的任务坑得焦头烂额,我建议你先按这套思路操作一遍:打开Spark UI或者Hadoop Application页面,看看掉队task的分布和输入数据量,判断到底是节点问题还是数据问题,然后再决定推测执行的开关和参数。这样做一次,你对推测执行的理解会比看十篇文档都深刻。