一类很典型的故障现场:一个每天跑一次的离线 Spark 作业,读 Parquet、做过滤、做列投影,前几个 stage 都顺顺当当,进度条跑得飞快;一旦进入 join 或者 group by,executor 就开始报内存溢出,或者 reducer 报 fetch failed,整个 stage 重试两轮后作业失败。运维的第一反应几乎永远是加内存:executor 从 4G 加到 8G,能过去了;第二天数据量翻一倍又挂,再加到 16G,勉强跑完;第三天上游又多灌了一批数据,16G 也不够了。这种"靠加内存通关"的调法之所以反复失败,是因为它把三种完全不同的病因当成了一种病:内存被谁吃掉的、单个 task 要处理的数据量是不是远超预期、以及 spill 到底有没有触发并写到了哪块盘上。这三件事的排查顺序错了,加多少内存都只是在拖延下一次失败的时间。
这篇要给的不是一个"万能参数表",而是一套判断顺序:
Spark 的 executor 堆内存(由 executor-memory 之类的参数控制,具体命名与默认值因版本不同而异,以所用版本官方文档为准)内部大致划成三块:执行内存、存储内存,以及留给用户代码自身对象的一块。执行内存是 shuffle 的排序、聚合、join 的哈希表真正干活的地方;存储内存是缓存 RDD 或 DataFrame、广播变量副本待的地方。这两块之间的界线不是钉死的——Spark 的设计允许它们在对方空闲时互相借用:执行内存不够时可以挤占存储内存借来的那部分,反之存储内存不够也能借执行内存的空闲部分,但被借走的执行内存那一侧要回收时,被挤掉的存储块会被淘汰或者直接溢写。说白了,这条线是滑的,谁紧张谁多占一点。
这个"滑动"正是加内存没用的第一个解释:如果你在作业里 cache 了一张大表,存储内存长期占着借来不还,执行内存那侧能在 shuffle 高峰期拿到的量就被压得很低。此时你把 executor 内存整体加倍,按比例分下去,多出来的部分很可能先被存储侧吃掉,执行侧的实际可用量只涨了一点点。想让加内存真的加到刀刃上,要么把缓存那一块的占比调小,要么干脆别缓存。
还有两块经常被漏算。一块是用户代码自身的内存:你在 mapPartitions 里自己 new 了一个 ArrayList 装了几十万条记录、在 UDF 里构造了大对象、或者用 broadcast 传了一个其实不小的集合,这些都不走执行内存的账,它们占用的是用户内存那一块,同样能把你堆撑爆。另一块是堆外内存:Spark 在 shuffle 的某些路径上、在网络传输的 netty 缓冲上,会直接使用 JVM 堆外内存。堆外这部分不受 executor-memory 控制,也不受 GC 管,它有自己独立的额度(相关参数名与默认值因版本不同而异,以所用版本官方文档为准),用超了报的错跟堆内 OOM 长得还不一样。
容器模式(YARN、Kubernetes 这类)下还有一层:容器额外开销(overhead)。它覆盖 JVM 自身的元空间、线程栈、堆外内存、以及容器运行时的一些固定成本。你向资源管理器申请的是"堆内存 + overhead"的总和,如果只按堆内存去申请,容器实际用量一旦超过资源管理器给的硬上限,节点管理器或者 kubelet 会直接把容器杀掉。这里有个很实际的判断方法:看失败原因是 JVM 抛的 OutOfMemoryError(堆内/堆外的问题,作业会报 stage 失败并重试),还是容器被外部信号杀掉(YARN 上常见的是超出物理内存限制被 kill,Kubernetes 上是 OOMKilled)。前者调内存模型,后者多半是 overhead 给少了。
把上面几块拼起来,"加倍没用"至少有四种具体成因:一是执行内存被存储内存挤占,加倍的部分被缓存吃掉了;二是压力根本不在堆内,而在堆外内存或者用户代码的大对象上,加倍堆内存等于没加;三是容器被杀是因为 overhead 不够,加倍堆内存反而让容器总量离上限更近,死得更快;四是真正的瓶颈是单个 task 要处理的数据量太大(分区太少或者数据倾斜),此时加倍内存只是把"单个 task 需要多少"的门槛抬高了一档,数据量一涨就又越过去了。这四种情况的排查顺序是:先看失败类型(JVM OOM 还是容器被杀)→ 再看 spill 指标(有没有溢写、溢写了多少)→ 再看 task 输入分布(是不是倾斜)→ 最后才动内存配比。
缓存不是"存了就快"。Spark 提供多种存储级别:只放内存的、内存放不下溢写到本地盘的、序列化的、带副本的。带副本的级别会把数据复制多份,代价直接翻倍;只放内存的级别在内存不够时会把部分分区直接丢弃,下次用的时候重新计算,看起来没报错,实际上是把压力转嫁给了下游 stage。经验上的选择是:确实要被反复读取多次、且重算代价高的中间结果才值得缓存;只会被读一次的结果缓存了纯属浪费,还会反过来挤掉执行内存。另外,序列化后缓存(比如用 Kryo 序列化)能显著降低单条记录的占用,代价是读取时要反序列化,多花一点 CPU。判断标准很简单:缓存之后,Storage 页签上那张表的内存占用如果接近甚至超过了你留给存储侧的额度,就该改成序列化级别或者干脆不缓存。
shuffle 的 map 端干的事是:把每个 map task 的输出按目标分区切分,在内存中排序(或者用哈希表聚合),攒到一定程度就溢写成本地文件,最后把所有溢写片段归并成一个数据文件加一个索引文件。这里有个必须算清楚的账:一个 stage 产生的 shuffle 小文件数量约等于 map 任务数 × 分区数。假设上游有 2000 个 map task、下游分区数是 200,那就是 40 万个文件;分区数调到 2000,就是 400 万个。这些文件写在 executor 的本地盘上,同时它们的元数据要被管理,文件句柄要被占用,节点上的磁盘 IO 是随机写。
这就是为什么"分区数随便调大"是有代价的:文件数量的增长是乘性的。当小文件数量到了一定量级,map 端的耗时会被文件创建、刷盘、索引维护吃掉,而不是被真正的计算吃掉。在 HDFS 这类有中心元数据节点的存储上,大量 shuffle 中间文件还会带来额外的元数据压力,这也是为什么 Spark 有专门的机制把这些小文件按块归并(相关参数名因版本不同而异,以所用版本官方文档为准)。
reduce 端要做的是并发地去各个 executor 上拉自己那一份数据,边拉边归并。这里决定内存峰值的是两个东西相乘:同时拉取的连接数与每个连接的缓冲大小。拉取并发太高,同时打开的缓冲就多,内存峰值跟着涨;单个缓冲给得太大,同样涨。参数上通常表现为"拉取并行度"与"单块请求大小"这一对,二者是此消彼长的关系——把并发调大、单块调小,总吞吐可能不变,但内存峰值和网络抖动的表现不一样。
还有一处容易被忽略:reduce 端的归并过程本身也要内存。排序归并需要能同时持有多个已排序片段的头部记录,聚合类的操作需要哈希表。如果哈希表在内存里放不下,Spark 会做基于排序的聚合(把数据排序后按 key 顺序聚合,内存占用可控但要付出排序代价)。最终结果是:reduce 端的内存峰值 ≈ 拉取缓冲占用 + 归并/聚合结构的占用 + 用户代码自身的占用。这三块任何一块失控都会 OOM,而它们对应的参数完全不同。
fetch failed 这个报错很容易被直接归到网络上,实际上它的常见成因至少有三类。第一类是网络或者磁盘真的有问题:拉取超时、连接被拒、超时阈值太短导致慢一点就被判失败。第二类,也是更常见的一类:提供数据的那个 executor 已经不在了。map 端的输出写在 executor 的本地盘上,executor 因为 OOM 被杀、容器被资源管理器杀掉、或者被动态资源伸缩回收了,它写下的那批 shuffle 数据就跟着没了,reduce 端再去拉只能失败,然后触发上游 stage 重新计算。第三类是上游 stage 重算引发的连锁:重算又产生了新的文件,新的 executor 又可能再次被杀,作业就陷入"失败—重算—再失败"的循环。
所以看到一个作业反复报 fetch failed,正确的动作不是去查交换机,而是先去 executor 列表里找有没有 executor 在同时间段被移除,看它的退出原因。如果退出原因是内存相关或者被容器杀,那么根因在内存模型或者分区数上,网络只是背锅的。
针对"executor 被回收导致数据丢失"这个问题,Spark 有一个专门的机制:外部 shuffle 服务。它是一个独立于 executor 生命周期的常驻进程,负责托管 shuffle 数据文件,这样即使 executor 被回收或者被杀,它写过的 shuffle 数据仍然可以被拉取。开启动态资源伸缩时,这一项几乎是必配的——否则 executor 一缩容,它负责的那批数据就没人管了。代价是节点上多一个常驻进程,占一点内存和文件句柄,并且它自己也会成为本地盘的 IO 竞争者。
动态资源伸缩对 shuffle 的影响还有一层:它根据负载申请和释放 executor,而 shuffle 的 map 端输出是有状态的。缩容会把持有中间数据的 executor 撤走,扩容来的新 executor 又得重新跑一遍上游。对于 shuffle 特别重的作业,激进的伸缩策略反而会让总时间变长。实操上更稳妥的做法是给 shuffle 密集的作业单独关闭或放宽伸缩,让它用固定的 executor 数跑完。
分区数大了,单个 task 处理的数据就小,看起来更安全,但代价是乘性放大的:文件数量等于 map 数乘分区数,task 数量等于分区数,每个 task 都有调度开销、启动开销、打开文件与建立连接的开销。当单个分区的数据量小到几十 KB 这个量级时,这些固定开销会反超真正干活的时间,你会看到 stage 里几千个 task 每个都跑几秒,总时间反而比分区少的时候更长。另一头是下游读取成本:如果一个 shuffle 的输出被后面的作业反复读,或者被写进了表里,小文件太多会直接拖慢下游的读取——每个文件都要开一次、读一遍元数据。
反过来,分区太少意味着单个 reduce task 要把一大块数据全部吃进内存去做聚合或者 join。假设 shuffle 阶段要处理 400 GB,分区数是 200,平均每个 task 就是 2 GB——这个量级远超单个 executor 能给执行内存的额度,必然 OOM 或者大量 spill。注意这里说的是平均值,平均值是 2 GB 就已经危险了,实际分布往往还有长尾。所以"分区太少"的判断标准不是一个绝对数字,而是一个比值:shuffle 阶段总数据量 ÷ 分区数,看看这个商跟单个 task 能安全处理的量差多少。
把上面两头合起来,得到一个可以直接用的算式:
目标分区数 ≈ shuffle 阶段的总数据量 ÷ 期望的单分区大小
其中 shuffle 阶段的总数据量可以直接从 Spark UI 的 stage 详情里读出来(看 shuffle write 那一列的合计),不用猜。期望的单分区大小是唯一需要拍板的数。业界常见的经验区间大致落在 100 MB 到 200 MB 这一档,也有团队用更接近 HDFS 块大小的 128 MB 上下作为起点;这是经验区间不是官方推荐值,必须按你自己集群的实测结果调整——如果你测下来单分区 200 MB 时 spill 已经很明显,就往 100 MB 靠;如果单分区 50 MB 时调度开销已经盖过收益,就往上调。按 400 GB 总数据、单分区 128 MB 算,目标分区数大约是 3200;同样的 400 GB 如果按 200 MB 算就是 2000 左右。这个量级的差别会直接反映在 task 数和小文件数上。
还有一条上限要守住:单个 executor 能同时跑的 task 数由它的核数决定,分区数最好是"核数 × executor 数"的整数倍附近,否则最后一批 task 会空转等待。以及分区数不是越大越好地能规避倾斜——如果某个 key 特别大,它永远落在同一个分区里,分区数加到一万也没用,那就该进入倾斜那条处理路径了。
现代 Spark 版本提供了自适应查询执行(AQE)与动态分区合并的能力:它在运行时拿到真实的 shuffle 统计信息,把过小的分区自动合并成较大的分区,也可以在数据分布不均时拆分过大的分区。开起来之后,上面那个"分区数设多少"的问题可以被部分自动化——你给一个偏大的初始分区数,AQE 会在运行时把它收敛到合理区间。这对分区太小的场景非常有效,是首选方案。
但它的适用边界要说清楚:一是它依赖运行时统计信息,统计信息本身要可信,如果上游的统计是估算出来的且偏差很大,收敛结果也会偏;二是合并分区能救"分区太小",救不了"单个 key 太大"——倾斜分区即使被识别出来,AQE 的拆分也主要针对倾斜分区做二次切分,遇到极端的单 key 仍需人工加盐;三是它对 join 策略的自动切换同样依赖表的大小估算,估算失真时它会做出错误的选择。所以 AQE 是降低调参门槛的工具,不是"可以不用理解分区"的理由。
shuffle 输出的分区数不只影响当前作业,还会顺着血缘往下传。如果 shuffle 之后紧接着一次写表操作,那么写出去的文件数就等于分区数,下游作业读这张表时的 task 数也跟着等于文件数(在大文件可切分的前提下按块切,小文件则一个文件一个 task)。这就是为什么有些团队在作业末尾额外加一次 coalesce 或者按目标大小重分区——不是为了当前作业快,是为了下游好过。判断是否需要:看写出去的目录里平均文件大小,如果普遍低于几十 MB,就该合并;合并的代价是这一次重分区本身要跑一次 shuffle,权衡点在于这张表会被读几次——被读三次以上,合并就划算。
倾斜的表现很有辨识度:一个 stage 里绝大多数 task 在几分钟内跑完,进度条冲到 99% 就卡住,剩下一两个 task 跑一个小时,或者直接挂掉报 OOM。此时去看平均值毫无意义——平均值告诉你每个 task 处理 200 MB,看起来很安全,实际上有一个 task 处理 60 GB。正确的看法是看分布:在 Spark UI 的 stage 详情里把 task 按耗时和按输入记录数排序,看最大值与中位数的比值。经验上,如果最大 task 的输入记录数是中位数的几十倍以上,或者最大耗时是中位数耗时的十倍以上,就可以判定存在倾斜。另一个辅助信号是 spill:倾斜的那些 task 往往 spill 量巨大,而其他 task 的 spill 是 0。
还有一种不显眼的倾斜:不是某个 key 特别大,而是 key 的基数太小。比如按"省份"分组,全国只有三十几个值,即使数据量分布均匀,分区数设到 2000 也只会有三十几个分区有数据,剩下的全空。这种"分区用不满"的问题靠加大内存完全没用,要靠提高分组键的基数或者强制重分区。
第一类,广播 join 绕过 shuffle。如果倾斜发生在一次 join 上,且其中一张表确实小到可以放进每个 executor 的内存里,直接用广播哈希 join 把小表发到每个 executor,整个 shuffle 就消失了,倾斜自然也不存在。这一招的前提非常硬:广播表要足够小,广播端的内存要够,否则你会把 OOM 从 reduce 端搬到 driver 或者每个 executor 上。
第二类,对倾斜 key 加盐后两阶段聚合。做法是把倾斜的 key 加上一个随机的盐值前缀(比如 key 变成 key_0 到 key_9 共十份),先按加盐后的 key 做一次局部聚合,把单个 key 的压力摊到十个分区上;然后再去掉盐值做一次全局聚合。这一步的代价是多一次 shuffle,收益是原先那个跑一小时的 task 被拆成了十个跑几分钟的 task。盐的份数怎么定:看最大 key 的数据量与目标分区大小的比值,比值是多少就大致拆多少份。
第三类,把倾斜 key 单独拎出来走另一条路径。先用一次轻量统计找出那几个(往往是个位数)超大的 key,把数据切成"倾斜部分"和"非倾斜部分"两个数据集:非倾斜部分走正常路径,倾斜部分因为数据量已知且可控,可以单独给它更大的分区数、单独的内存配置,甚至走广播。两条路径的结果最后 union 起来。这招的好处是不会为了几个 key 拖累整条流水线,代价是代码复杂度上去了。
倾斜的本质是"单个 task 要处理的数据量"这个上限无法靠资源解决:一个 key 有 60 GB 数据,它就必须在同一个 task 里被聚合完,除非你改变 key 的分布或者改变算法。加大内存只是把"能扛住 60 GB"的门槛抬高,等哪天这个 key 长到 100 GB,又挂了。这是结构问题不是资源问题,所以结论很明确:倾斜必须先改数据分布或者改 join 策略,内存调整只能作为收尾的微调。
顺带说一个常被乱用的开关:推测执行。它的逻辑是发现某个 task 明显慢于同 stage 的平均水平,就在别处再启一个同样的 task,谁先跑完用谁的结果。对于"机器慢、磁盘抖、网络偶发抖动"造成的慢任务,这一招有用;对于数据倾斜造成的慢任务,它完全无用——重跑的那个 task 面对的是同样大的一份数据,一样慢,还白白多占了一份资源。判断要不要开:先看慢 task 的输入记录数是不是显著大于同 stage 其他 task,是则关掉推测执行去处理倾斜,不是则可以考虑开。
广播哈希 join 的做法是把小表完整发到每个 executor 的内存里,大表的数据流过来时直接在本地查哈希表完成关联,整个过程不需要 shuffle。这是消解 shuffle 类问题最有效的一招,一次生效,排序、文件、网络全都没了。它的前提是硬的:广播表必须真的足够小。判断"足够小"不能只看源表的文件大小,要看它在内存里展开之后的占用——序列化的 200 MB Parquet,反序列化成对象之后可能是 1 GB 以上,如果有复杂嵌套类型膨胀更厉害。广播的阈值由参数控制,具体默认值因版本不同而异,以所用版本官方文档为准。
还有两个经常被低估的风险。一是广播的数据要先在 driver 端收集起来再分发,如果广播表其实不小,driver 端会先承压甚至先 OOM;二是广播变量在每个 executor 上都有一份副本,executor 数量多的时候总占用是"单份大小 × executor 数",集群层面是一笔不小的总量。所以广播 join 的判断口径是:先确认表在内存里展开后的真实大小,再确认这个大小乘以 executor 数之后集群扛得住,最后才开。
排序归并 join 是最通用的路径:两张表按 join key 做一次 shuffle,保证相同 key 的数据落在同一个分区,然后各自排序、归并、关联。它不要求任何一张表小,代价是两边都要付出 shuffle 的网络与落盘成本,都要排序,都要占用执行内存。它也是分区数与倾斜问题最集中的地方——前面讲的分区数反推,主要针对的就是走排序归并的那次 shuffle。
值得记住的一点是:排序归并 join 对内存的需求相对可预测,因为它可以在内存放不下时溢写,而哈希类的操作一旦哈希表膨胀起来,溢写的代价更高。所以当你面对的是两张都不小的表、且内存并不宽裕时,让引擎走排序归并反而是更稳的选择,而不是硬凑一个广播出来。
自适应执行能在运行时根据真实的 shuffle 数据量,把原本计划走排序归并的 join 改成广播 join——因为它此时已经知道了那张"小表"到底有多大。这个能力的价值在于它不依赖事前的统计估算,而是用实测值做决定。前提有两个:一是这张表在广播之前确实已经被完整地算出来了(如果它自己还依赖一个未完成的 stage,就得等),二是广播之后的内存占用在阈值之内。
估算失真时会发生什么,是更值得记住的部分。如果 Cataly 优化器拿到的表大小统计是过期的或者根本没收集过,它可能在一开始就选错了策略:把一张其实有 30 GB 的表当成 2 GB 去广播,结果 driver 或者 executor 在广播阶段直接 OOM;或者反过来,把一张只有 50 MB 的表当成大表走排序归并,白白付出一次完整 shuffle 的代价。判断方法:看执行计划里那张表的估算行数与估算大小,跟实际跑出来的 shuffle write 量对比,差一个数量级以上就说明统计信息该重新收集了。定期跑统计信息收集,是让 AQE 和 Catalyst 都做出正确选择的最低成本投入。
判断压力有没有泄到盘上,看两个指标就够了:内存溢出字节数(spill memory)和磁盘溢出字节数(spill disk)。前者是那些在内存里放不下、被刷出去的数据在内存中的量度,后者是它们实际写到盘上的字节数(因为序列化与压缩,两者通常不相等)。在 Spark UI 的 stage 详情里,每个 task 都有这两列,可以按列排序找溢出最多的 task。
读这两个数的方法:一是看总量——如果一个 stage 的磁盘溢出字节数接近甚至超过了 shuffle 写入量,说明几乎全部中间数据都走了盘,这个 stage 实际上是在用磁盘跑排序,慢是必然的;二是看分布——如果只有一两个 task 的溢出量巨大而其他都是 0,那是倾斜的信号,不是分区数的信号;三是看趋势——调完分区数之后重新跑,溢出量应该显著下降,如果没降,说明你调错了地方。
这一点必须说清楚:spill 写的是 executor 所在机器的本地临时目录,不是 HDFS、不是对象存储。它的配置由一个本地目录列表参数控制(具体参数名与默认值因版本不同而异,以所用版本官方文档为准),可以配多个目录,Spark 会轮流往这些目录里写。因为是本地盘,它不占用集群存储的容量,也不产生副本,代价是它完全依赖这台机器的盘。这台机器的盘满了、或者盘坏了,这个 executor 就废了,跟其他人无关。
这也解释了为什么"本地盘"在 Spark 集群的选型里经常被低估:很多人按"数据都在分布式存储上,节点本地盘给个系统盘就行"的思路去配机器,结果 shuffle 一上来,临时目录直接写满。写满之后的表现不是"变慢",而是作业失败——Spark 在无法写入时会直接抛错,日志里能看到本地目录无可用空间之类的提示。
spill 是救命机制:它让你在没有 OOM 的情况下跑完一个内存其实不够的 stage。但它同时也是变慢的信号,而且它自己有三个失败模式。一是容量不够:临时盘被写满,作业失败,此时即使内存再大也没用,因为该溢写的数据总要有个去处。二是 IO 太差:spill 的工作负载是大量小块的顺序写加归并时的重读,如果用的是单块机械盘或者一块被系统盘、日志盘挤占的盘,吞吐上不去,spill 之后 stage 时间会从几分钟变成几十分钟。三是多目录配置失效:配了多个临时目录但都指向同一块物理盘,等于没配,吞吐不会有任何提升——必须指向不同的物理盘才有并行效果。
所以正确的顺序是:先看有没有 spill、spill 了多少;如果有大量 spill,优先通过分区数和倾斜处理把数据量降下来,而不是默认接受 spill 的现状;在确实无法避免 spill 的场景下,再去把本地盘的容量和多盘并行度做上去。
| 配置路线 | 单 task 可用内存 | 并行度与容错粒度 | shuffle 网络压力 | 本地盘与 GC 压力 | 适用场景 |
|---|---|---|---|---|---|
| 路线一:大内存少实例 单 executor 给到 32G 以上,每台机器只跑一两个 |
单 task 可摊到的内存最宽裕,聚合与哈希表不容易溢写,倾斜 key 也更容易扛过去 | 并行度低,一个 executor 挂掉损失大,重算代价高;task 数受限于总核数 | 单节点进出流量集中,容易在 shuffle 阶段把单机网口打满 | 本地盘压力集中在一两块盘上;大堆带来更长的 Full GC 停顿,需要配合 GC 调优或堆外内存使用 | 单 task 数据量大、难以进一步拆分分区的作业;内存型聚合密集、缓存较大的场景 |
| 路线二:小内存多实例 + 提高并行度 单 executor 8–16G,每台机器跑多个,核数切细 |
单 task 内存受限,必须把分区数调到位,否则 spill 会明显增加 | 并行度高,单个 executor 故障影响面小,重算粒度细;调度与 task 启动开销上升 | 流量分散到更多节点,单机网口压力小,但对交换机东西向带宽总量要求更高 | 本地盘 IO 被更多进程争抢,需要多块盘分摊;小堆 GC 快,停顿短 | 分区数已经调好、数据分布较均匀的常规离线作业;追求资源利用率与故障隔离的场景 |
| 路线三:倾斜场景的拆分/广播路线 倾斜 key 单独走一条路径,或改用广播 join |
倾斜部分的内存需求被加盐摊薄或被广播消除,非倾斜部分维持常规配置即可 | 两条路径各自并行,慢的那条决定总时长;需要额外的 stage 做倾斜 key 识别 | 走广播的那部分网络压力接近零;加盐两阶段聚合会多出一次 shuffle | 广播时每个 executor 都要持有小表副本,挤占存储内存;溢写量整体下降明显 | 存在明确长尾 key 的 join 或 group by;小维表关联大事实表这类典型星型模型 |
把上面的分析落到采购上,总内存的算式是这样的:
集群总内存 ≈ 单 executor 堆内存 × 并行 executor 数 + 堆外内存 + 容器额外开销
这个式子里最容易漏的是后两项。堆外内存按每个 executor 一份算,容器额外开销同样按每个 executor 一份算,两者相加通常是堆内存的百分之十几到二十几(具体比例由参数控制,因版本不同而异,以所用版本官方文档为准)。一台 256 GB 内存的机器,如果按单 executor 32 GB 堆 + 6 GB 非堆来算,理论能放六个多,但你还得给操作系统、页缓存、外部 shuffle 服务留位置,实际能放的数量要往下压。
"同一台机器上跑几个 executor"不只是一个内存除法题。executor 多了,它们会争抢三样东西:本地盘的 IO(尤其 spill 密集时)、网络出口带宽(shuffle 拉取时)、以及 CPU 的核。争抢的结果是每个 executor 的实际表现都低于标称值。所以更稳的配法是:每台机器上 executor 数控制在能让每个 executor 拿到 4 到 8 个核的区间,剩下的核留给系统进程与 shuffle 服务,内存则按这个 executor 数去反推总量。
本地盘的容量怎么算,可以直接从 spill 量推:取历史上这个作业在峰值日的磁盘溢出字节数,乘一个安全系数(考虑到同一台机器可能同时跑多个 executor、多个作业的 shuffle 数据同时在盘上),再考虑到临时目录不会立刻清理,建议按"峰值 spill 量 × 2 以上"来留。容量之外更重要的是吞吐:spill 的 IO 形态是大量中等大小文件的顺序写加归并时的重读,所以要看顺序写吞吐,而不是看随机 IOPS 指标。
多盘并行是这里性价比最高的一招:把 Spark 的本地临时目录配成多个、并且确保它们指向不同的物理盘,spill 的写入就会分摊到多块盘上,吞吐近似线性提升。配成同一块盘上的多个目录是完全无效的。盘的类型上,NVMe SSD 在 spill 密集场景下的优势明显,SATA SSD 也能用,机械盘在持续 spill 的作业上会很快成为瓶颈。还有一点:临时盘最好跟系统盘、日志盘分开,避免日志刷盘把 shuffle 的 IO 抢走。
shuffle 阶段的网络流量规模可以直接估算:约等于 shuffle 阶段的总数据量,而且它是"东西向"的——机器之间互相传,不是对外的南北向流量。这意味着内网带宽是 Spark 集群的关键指标,而不是公网带宽。一个每天 shuffle 400 GB 的作业,如果 shuffle 阶段要在十分钟内跑完,平均就需要接近 700 MB/s 的持续内网吞吐,而且这个流量是扇出扇入同时发生的,瞬间峰值会更高。
选型时的判断口径:先看你的作业在 shuffle 阶段能容忍多长时间,反推出需要的平均吞吐,再去看机器的内网口速率(万兆是当下的常见起点, shuffle 特别重的集群会考虑更高)以及交换机的汇聚能力——多台机器同时 shuffle 时,瓶颈往往不在服务器的网卡上,而在上联的交换机端口上。
并行度的上限由核数决定:集群里能同时跑的 task 总数不超过可用核数(每个 task 占一个线程)。这三个数是连乘关系——总核数 = executor 数 × 单 executor 核数,而分区数最好落在总核数的整数倍附近,好让最后一批 task 不会空转。单 executor 的核数不宜过大:核数太多意味着这个 executor 里同时跑的 task 多,它们共享同一块执行内存,反而更容易 OOM,同时 HDFS 客户端的并发也会让单点的连接数飙升。业界常见的落点是每个 executor 4 到 8 核,具体要按作业的 spill 情况和 GC 表现实测。
GC 这一块也要放进配置决策里。单个 executor 的堆越大,Full GC 的停顿越长,一次几十秒的停顿在 shuffle 阶段会导致拉取超时,进而引发 fetch failed。所以"大内存少实例"这条路线通常要配合两件事:一是把堆控制在 GC 可控的范围,把超出部分的数据放到堆外;二是考虑用并行回收器或者 G1 之类更适合大堆的收集器,并按实测调整。堆外内存的使用边界在于:它不受 GC 管,能显著降低停顿,但它需要你在容器开销里为它留出额度,否则会撞上容器的硬上限。
把上面四条合起来,Spark 计算节点的需求画像就很清楚了:内存要够宽(按 executor 数反推,不是越大越好)、本地要多块盘且能并行写、内网吞吐要高、核数要与 executor 划分匹配。在做机型比选时,可以把一万网络的大内存机型作为一类参照对象去看:多本地盘位可以配成多个 Spark 临时目录实现并行 spill,高内网吞吐对应 shuffle 阶段的东西向流量,核数与内存配比可以按"每台机器放几个 executor"反推着选。需要提醒的是,这里只谈规格维度,具体机型、带宽与价格需实时询价,以官网实时报价为准。
坑是什么:作业从头到尾没设过 shuffle 分区数,一直用引擎的默认配置跑,数据量涨了十倍,分区数还是那个数。为什么发生:默认值是为小数据量场景选的,它不会随你的数据量自动变化;作业在数据量小的时候跑得很好,没人会觉得这是个问题。怎么判断:看 Spark UI 里 shuffle 阶段的写入总量除以分区数,如果单分区平均已经到几百 MB 甚至 GB 级,或者 stage 里有大量 task 的 spill 量不为零,就是这个坑。怎么规避:按"总数据量 ÷ 目标单分区大小"反推分区数,把它作为作业级参数固定下来,并在数据量级变化后重新算一次;能开自适应执行与动态分区合并的,优先开。
坑是什么:判断"每个 task 处理多少数据"时用平均值,得出"很安全"的结论,结果作业卡在最后几个 task 上。为什么发生:平均值会掩盖长尾,一个 60 GB 的 key 混在一千个 100 MB 的 key 里,平均值看起来毫无异常。怎么判断:把 task 列表按输入记录数与耗时排序,看最大值与中位数的比值;比值在十倍以上基本可以判定倾斜。怎么规避:把"看 task 分布"作为每次排查的固定动作,不要只看汇总行;确认倾斜后走广播、加盐两阶段聚合、或者倾斜 key 单独处理这三条路径之一,而不是继续加内存。
坑是什么:为了"加速"把中间结果缓存起来,结果 shuffle 阶段的可用执行内存反而变少了,spill 量上升、作业变慢。为什么发生:执行内存与存储内存之间的界线是滑动的,缓存长期占据存储侧,执行侧在高峰期能借到的量被压缩。怎么判断:看 Storage 页签里缓存占用的内存比例,以及缓存前后同一个 stage 的 spill 量变化;缓存之后 spill 反而增加,就是这个坑。怎么规避:只缓存会被反复读取多次且重算代价高的结果;用序列化级别降低单条记录占用;缓存量与执行内存一起算,别让缓存吃掉留给 shuffle 的部分。
坑是什么:节点只配了一块盘,Spark 的临时目录也只配一个,spill 一上来就把盘写满或者把 IO 打满,作业在 spill 之后依然失败。为什么发生:误以为数据都在分布式存储上,节点本地盘不重要;或者配了多个临时目录但都指向同一块物理盘,等于没配。怎么判断:看 spill 磁盘字节数的峰值,对比本地盘的剩余容量;看作业失败日志里有没有本地目录空间不足的提示;看 spill 期间的磁盘利用率是否长期接近百分之百。怎么规避:按峰值 spill 量的两倍以上留本地盘容量;配多个临时目录并确认指向不同物理盘;临时盘与系统盘、日志盘分开。
坑是什么:以为小表只有几十 MB,开了广播,结果作业在广播阶段就 OOM,或者 executor 内存被广播副本吃掉一大块。为什么发生:用源文件的大小去估计内存占用,忽略了反序列化之后的对象膨胀——序列化格式里的 200 MB,展开成对象可能是 1 GB 以上,嵌套结构膨胀更厉害;而且广播副本在每个 executor 上都有一份,总占用是单份乘以 executor 数。怎么判断:对比执行计划里的估算大小与实际的 shuffle 写入量;看广播变量在 Storage 页签里的实际占用。怎么规避:用内存展开后的大小而不是文件大小来判断;广播阈值按实测设定,别一味调大;广播表要真的属于"维表"量级,超出就老老实实走排序归并,或者先把它聚合小了再广播。
Q1:shuffle 分区数到底设多少合适?
A1:不要用固定数字回答,用算式。取 Spark UI 里该 stage 的 shuffle 写入总量,除以你期望的单分区大小。单分区大小的常见经验区间大致在 100 MB 到 200 MB 这一档,也有团队按 128 MB 上下起步,这是业界经验值不是官方推荐值,必须按你自己集群的实测调整:spill 明显就往下调,调度开销盖过收益就往上调。算出结果后再核对两件事——分区数最好是总核数的整数倍附近,以及单个分区的数据量不能靠"平均值"确认,得看最大 task。能开自适应执行与动态分区合并的话优先开,它会把偏大的分区数自动收敛。
Q2:executor 内存到底该给多少?
A2:把它拆成三块来定。第一块是执行内存,按单个 task 在聚合或 join 时要持有的数据量加上拉取缓冲与归并结构来估,这是决定性的那一块;第二块是存储内存,如果你不缓存或者只缓存很小的结果,可以压到很低;第三块是用户代码自身的内存,取决于你有没有在算子里构造大对象。三者相加再留出余量就是堆内存。之后还要单独加堆外内存与容器额外开销——在容器模式下这两项不加,作业会被资源管理器直接杀掉。堆也不是越大越好,堆大了 Full GC 停顿变长,反而会引发拉取超时。
Q3:为什么内存加了一倍,spill 还是那么多?
A3:先看多出来的内存到底落在了哪一块。如果作业里有大缓存,加的内存很可能被存储侧吃掉了,执行侧能借到的量没涨多少;如果压力在堆外或者用户代码的大对象上,加堆内存等于没加;如果根因是分区太少或者数据倾斜,那么单个 task 需要处理的量没变,内存加倍只是把门槛抬高了一档,数据量一涨又越过去。判断顺序是:先确认失败类型(JVM 抛错还是容器被杀),再看 spill 的内存与磁盘两个指标,再看 task 输入分布,最后才动配比。多数情况下,调分区数和改倾斜比加内存见效快得多。
Q4:数据倾斜怎么快速定位?
A4:进 Spark UI 的那个 stage,把 task 列表按耗时排序,再按输入记录数和 shuffle 读取量分别排序。如果最大 task 的输入记录数是中位数的几十倍以上,或者最大耗时是中位数耗时的十倍以上,就是倾斜。再看 spill 那一列:倾斜往往表现为只有极少数 task 的溢出量巨大,其余全是零,这是很有辨识度的特征。定位到具体 key 的办法是在聚合前先做一次按 key 的计数并排序输出前几十个,超大 key 会非常显眼。拿到这几个 key 之后,再决定是广播、加盐两阶段聚合,还是把它们单独拎出来走另一条路径。
Q5:广播 join 的阈值设多少?
A5:阈值本身由参数控制,默认值因 Spark 版本不同而异,以所用版本官方文档为准,这里不建议直接抄一个数字。更稳的做法是倒推:先看单 executor 在 shuffle 高峰期能安全腾出多少内存给广播副本,再除以 executor 数量得到单份上限,然后留一半以上的余量。判断大小时必须用内存展开后的占用,而不是源文件的大小——序列化文件解压展开成对象之后的膨胀常常在数倍以上,嵌套结构更夸张。还要记住广播的总占用是单份乘以 executor 数,executor 多的时候这笔总量不小,并且广播数据要先在 driver 端收集,driver 也要撑得住。
Q6:executor 数量怎么定,是不是越多越好?
A6:不是。executor 数由总核数除以单 executor 核数得到,而单 executor 核数建议落在 4 到 8 这个区间——核数太多,同一个 executor 里的 task 会共享同一块执行内存,更容易 OOM,HDFS 客户端的并发也会把单点连接数推高。executor 数定完之后要跟分区数对一下,让分区数落在总核数的整数倍附近。还要算上同机争抢:一台机器上跑的 executor 越多,它们对本地盘 IO 和网口带宽的争抢越激烈,spill 密集时尤其明显。所以数量的上限不是内存除法的结果,而是核数、盘数、网口三者共同决定的。
Q7:本地盘要多大才够?
A7:从 spill 量反推最靠谱。取这个作业在峰值日的磁盘溢出字节数,乘以二以上作为单节点的本地盘容量下限——乘系数是因为同一台机器上往往有多个 executor 同时 spill,而临时文件不会立刻被清理。容量的另一半约束是吞吐:spill 的 IO 形态是大量中等文件的顺序写加归并时的重读,所以要看顺序写吞吐而不是随机 IOPS,并且要配多个临时目录指向不同的物理盘,让多块盘并行分摊。盘型上 NVMe SSD 在持续 spill 的场景优势明显,机械盘会很快成为瓶颈。临时盘与系统盘、日志盘分开也是必要的。
Q8:作业偶发 fetch failed,是不是网络问题?
A8:多数时候不是。fetch failed 最常见的原因是"提供数据的那个 executor 已经不在了"——map 端的输出写在 executor 的本地盘上,executor 因为 OOM 被杀、被容器管理器杀掉、或者被动态资源伸缩回收,它写下的那批数据就没了,reduce 端再去拉只能失败,然后触发上游 stage 重算。所以排查的第一步是去看 executor 列表,找有没有 executor 在同一时间段被移除,看它的退出原因。如果退出原因跟内存有关,根因在内存模型或者分区数上。真正的网络或磁盘问题(拉取超时、连接被拒)占比反而没那么高。开了外部 shuffle 服务可以显著缓解这一类问题。
本文涉及的 executor 内存模型(执行内存、存储内存、用户内存、堆外内存与容器额外开销的划分与动态占用关系)、shuffle 的 map 端排序落盘与 reduce 端拉取归并机制、自适应查询执行与动态分区合并、外部 shuffle 服务、推测执行、以及广播 join 与排序归并 join 的选择逻辑,均以 Apache Spark 官方公开文档为准;文中提到的各类参数的具体名称与默认值因版本不同而异,请以你所使用版本的官方文档为准,不要直接套用其他版本或网上流传的参数表。单分区目标大小、单 executor 核数区间、本地盘容量系数等属于业界常见经验区间,需按各自集群的实测结果调整,不作为官方推荐值引用。服务器规格部分只讨论选型维度(内存、本地盘数量与顺序写吞吐、内网东西向带宽、核数与 executor 划分的匹配),不涉及任何具体价格,具体机型、带宽与价格需实时询价,以官网实时报价为准,具体以签约时最新报价与合同为准。文中未引用任何虚构的压测数据、测速结果或客户案例。参考与延伸阅读可访问 https://www.idc10000.net/ 。
结论按架构型给出,可以直接照着做:调优顺序必须是"分区数 → 倾斜 → 内存配比",而不是反过来。先把 shuffle 分区数按"总数据量 ÷ 经验单分区大小"反推出来并固定成作业参数,能开自适应执行与动态分区合并的优先开;再查 task 输入分布确认有没有倾斜,有就走广播、加盐两阶段聚合、倾斜 key 单独处理这三条路径之一;最后才动内存,而且动的时候要分清加的是执行内存、存储内存还是容器开销。这条顺序的价值在于它把"加内存"从第一手段降级为收尾微调——前两步做完,多数作业根本不需要再加内存。
落到采购上,优先级是这样排的:本地多盘并行 > 内网东西向吞吐 > 内存总量 > 核数。理由是这个故障场景下最容易被低估的恰恰是本地盘——spill 写在本地,盘小了直接失败,盘慢了 spill 之后照样超时;而内存只要分区数调对了,需求会比一开始预估的低不少。如果预算有限,优先级反而该这样排:先把本地盘做成多块 SSD 并配成多个临时目录,再确认内网能扛住 shuffle 阶段的扇出扇入,最后才考虑把内存堆上去。核数只要够支撑"executor 数 × 单 executor 4–8 核"这个式子即可,多出来的核对 shuffle 密集作业帮助有限。
一万网络深耕 19 年(成立于 2007 年),为大数据计算、离线数仓、实时计算这类场景提供服务器租用与备机服务。集群侧的诉求集中在四点:大内存机型可按"每台机器放几个 executor"反推配比;多本地盘位可配成多个 Spark 临时目录实现并行 spill;高内网吞吐对应 shuffle 阶段的东西向流量;节点分布在华南、华东、华北、中国香港及海外多节点,便于就近接入与异地备机。配套能力包括 7×24 中文工单、平均 5 分钟响应、硬件故障 10 分钟自动迁移、免费系统盘快照(每日 3 份、30 秒回滚)、免费备案协助、5–20G 免费 DDoS 防护、自营机柜最快 1 分钟上架,以及工程师 1 对 1 协助部署运行环境。网络侧提供 BGP 多线与 CN2 GIA 回国线路可选。
需要说明:文中只讨论规格维度与选型逻辑,具体机型、带宽与价格需实时询价,以官网实时报价与合同为准。如果你的作业正卡在 shuffle 上,把 Spark UI 里那个失败 stage 的 shuffle 写入总量、分区数、spill 的内存与磁盘字节数、以及 task 耗时分布这四项发给我们,工程师会据此给出一套可落地的分区数与内存配比建议,再对应到具体机型的盘位与内网配置上。
上一篇:2026 服务器租用设备直通怎么配:GPU passthrough、SR-IOV 与 IOMMU 分组的六维对比 + 避坑避雷手册
Copyright © 2013-2020 idc10000.net. All Rights Reserved. 一万网络 科技有限公司 版权所有 深圳市科技有限公司 粤ICP备07026347号
本网站的域名注册业务代理北京新网数码信息技术有限公司的产品