关于我们

质量为本、客户为根、勇于拼搏、务实创新

< 返回新闻公共列表

2026 Spark 作业总在 Shuffle 阶段变慢:中间数据落盘的磁盘该怎么配

发布时间:2026-09-28

reduce 阶段一卡就是十几分钟,问题往往不在 CPU

先看现象:一个跑了半年的 T+1 离线作业,某天开始超时。Spark UI 上 map 阶段 4 分多钟全部跑完,reduce 阶段进度条停在 199/200,剩下那一个 task 转了 20 分钟不动,然后 Executor 日志里开始刷 ExternalAppendOnlyMap 的 spill 记录,紧跟着一条 org.apache.spark.shuffle.FetchFailedException,Driver 把某个 Executor 标记成 lost,再让上游的 map 端重新算一遍。整个过程像死循环:重算 → 再落盘 → 再 fetch 失败 → 再重算。

再看监控:节点的 CPU 使用率不但没打满,反而从 70% 掉到 20% 上下;内存还剩三分之一;唯一异常的是 iostat 里落盘那块盘 %util 长期 100%,await 从平时的几毫秒涨到一两百毫秒,w/s 却高得离谱。这时候加 CPU、加内存,都不会让进度条动一格。

问题的根:Shuffle 是 Spark 里唯一一个「必须写本地磁盘」的环节。map 端的中间结果要落到 Executor 所在机器的本地盘上,reduce 端再跨网络把它们拉回来。磁盘一旦跟不上,CPU 就只能空转等 IO,表现出来就是 reduce 卡住、fetch 超时、Executor 被判定失联。换句话讲,Shuffle 阶段的瓶颈,十有八九是在本地盘的容量、顺序写吞吐和随机写能力上,而不是在计算上。

本篇的落点:把 shuffle 中间数据落盘这件事拆成三个可量化的变量——容量怎么算、顺序写要多少、随机写要多强——再对应到 NVMe / SATA SSD / HDD 三种盘该怎么选、几块盘、多大、配什么参数,最后落到一万网络上能直接下单的两组对照配置。

Shuffle 到底往磁盘上写了什么

要配盘,先得知道盘上到底躺着什么东西。以默认的 SortShuffleManager 为例(具体管理器类型与版本行为以官方文档对应版本为准),一次 shuffle 在磁盘上产生两类东西:

map 端:每个 task 一组 data 文件 + index 文件

每个 map task 会把自己输出的数据按目标分区排好序,写进一个本地的 data 文件,同时生成一个 index 文件记录每个 reduce 分区在这份 data 文件里的偏移和长度。也就是说,如果有 M 个 map task,磁盘上至少会有 M 份 data 文件和 M 份 index 文件,而每份 data 文件内部又被切成 R 段(R 就是 reduce 侧分区数,SQL 场景下由 spark.sql.shuffle.partitions 决定)。

这里有个经常被误解的点:就算内存再大,map 端的 shuffle 数据也是要落盘的。sort-based shuffle 的设计就是把排序后的结果写成本地文件供下游拉取,不是「内存够就不写」。内存在这里只决定排序缓冲区有多大、能少刷几次,不决定写不写。

写的过程中,spark.shuffle.file.buffer 控制每个输出文件的内存缓冲大小,缓冲攒够了才真正往下刷。这个值调大能减少系统调用次数、让单次写入更成块,代价是每个并发输出流占用的内存变多——一个 map task 同时要写 R 个分区段,缓冲给得太大,Executor 内存会被这些 buffer 吃掉一大块。

reduce 端:拉回来的数据在内存里聚合,装不下就 spill

reduce 端通过 spark.reducer.maxSizeInFlight 控制同一时刻从远端拉取的数据量上限,拉回来的数据先进内存做聚合(比如 group by、join 的 hash 表)。当哈希表涨到阈值放不下时,就会把它溢写到本地盘,这就是日志里那句 spill。之后再把多个 spill 文件归并,归并阶段同样可能二次落盘。

spark.reducer.maxSizeInFlight 调大,能让拉取更饱满、减少轮次,但会让 reduce 端同时驻留更多未处理数据,内存压力上升;调小则拉取变慢、等待变多。它和磁盘的关系在于:拉取速率上去了,聚合和 spill 的速率也要跟得上,否则数据堆在内存里反而更早触发溢写。

合起来的 IO 画像

  • map 端写:以顺序追加为主,单个分区段内是连续写,但因为要同时写 R 个段,实际落到盘上会带一点随机性,盘多、分区多的时候更明显。
  • reduce 端读:从本地盘上读 map 端写的文件,是大量小块的随机读,扇区大小常常只有几十 KB。
  • reduce 端写:spill 是典型的随机小 IO 写,而且是突发式的——一瞬间要写几百 MB。
  • 文件数量:M×R 这个量级决定了 inode 消耗和文件句柄压力,分区数设得夸张时,元数据操作本身就能拖垮一块机械盘。

所以 shuffle 对磁盘的要求可以概括成一句话:写多读多,大顺序写和小随机 IO 混在一起。这也意味着,决定一块盘能不能扛住 shuffle 的,依次是「容量够不够装」>「顺序写吞不吞吐」>「4K 随机写强不强」。顺序不够会拖慢 map 落盘,随机不够会让 reduce 端 spill 和 fetch 卡成幻灯片。

三种磁盘在 shuffle 场景下的真实差别(NVMe / SATA SSD / HDD)

参数表上的「顺序读写 MB/s」是最容易误导人的一项。shuffle 的痛苦点在混合负载,尤其是 reduce 端那波突发随机写。

HDD:顺序写尚可,随机 IO 直接崩

单块企业级机械盘的顺序写通常在 150–220MB/s 量级(行业参考,具体以盘型实测为准),看起来不算差。但它的随机 4K 写能力只有一百到几百 IOPS,每次 IO 还要付磁头寻道的毫秒级延迟。放到 shuffle 里会发生什么?map 阶段还能勉强跑,一旦进入 reduce 端的 spill 与归并,几十上百个并发线程各写各的小文件,磁头在盘片上来回摆动,实际吞吐掉到十几 MB/s 都不稀奇,await 直接上百毫秒。表现就是开头那个现象:map 快、reduce 死。

HDD 还有一个隐藏成本——它是集群里最容易触发「慢节点效应」的部件。一个 reduce 分区要等所有 map 端的数据都到位,只要有一台 Executor 的盘慢,整个 stage 就被这一个节点钉住。用 HDD 跑 shuffle,等于给每个作业装了一个随机位置的地雷。

SATA SSD:够用的起步线,但要小心掉速

SATA 接口 6Gbps 的天花板决定了这类盘的顺序写基本停在 400–550MB/s 区间(行业参考),随机 4K 写能到数万 IOPS,比 HDD 高两个数量级。对日跑批量在几百 GB 级别、并发 task 不夸张的集群,多块 SATA SSD 是性价比很扎实的落盘方案。

需要盯住两件事。一是掉速:不少 TLC 盘靠 SLC 缓存撑标称写入速度,缓存写满之后直写 TLC,速度会掉到原来的三分之一甚至更低,而 shuffle 恰恰是一次连续写几十上百 GB 的场景,正好撞在这个坎上。二是写寿命:按「每日 shuffle 落盘量 × 365 ÷ 单盘 TBW」粗算,如果一年就把标称写入量耗掉,盘会在保修期内提前进入只读或坏块状态。这两项都属于(预估)范畴,选型时以盘厂规格书与实测为准。

NVMe:把随机 IO 这个短板补齐

NVMe 走 PCIe 通道,队列深度从 SATA 的 32 提到数万级,单盘顺序写常见 1–3GB/s,随机 4K 写能到几十万 IOPS 量级(均为行业参考区间,实际以型号实测为准)。shuffle 场景最吃的两项——并发小写和突发 spill——正好是 NVMe 的强项。

更关键的是延迟一致性。SATA SSD 在持续写入压力下延迟曲线会抖,NVMe 在同样压力下多数能保持比较平稳的亚毫秒级响应,这直接决定了 reduce 阶段会不会出现个别 task 突然慢十倍。对一个每天要在固定时间窗内跑完的 T+1 作业来说,稳定的尾延迟比峰值带宽更有价值。

当然 NVMe 不是白拿的。单盘位价格高于同容量 SATA SSD(价差幅度随容量与市场波动,属(预估)范畴,以咨询为准),而且要注意 PCIe 通道数:一台机器上插满多块 NVMe,如果通道被 CPU 和网卡分走,实际能跑满的盘数是有上限的,选型时要看主板的通道分配,别买了四块盘结果只有两条通道的带宽。

容量怎么估:一个可以自己算的粗算口径

容量是最容易拍脑袋的一项,也是最容易出事的——盘写满不是「慢一点」,而是 Executor 直接被集群判死、作业重算。下面给一套能自己往上套数字的口径。

第一步:拿到单次作业的 shuffle 写盘量 W

最准的办法是看 Spark UI 的 Stage 详情页,Shuffle Write 那一列给出的就是 map 端实际写出去的字节数,把作业里所有 shuffle stage 的写盘量加起来得到 W。历史作业跑过一次就有数,比任何公式都准。

如果作业还没上线、拿不到 W,可以按输入量反推(属粗略估算):

  • 大表 join:shuffle 量约为参与 join 的输入量之和的 1–3 倍(预估),宽表、多列参与、带膨胀的 join 取上限;
  • 分组聚合:shuffle 量通常小于输入量,约为 0.3–1 倍(预估),取决于聚合前的字段投影和预聚合是否生效;
  • 去重 / distinct / 窗口函数:膨胀系数波动最大,1–4 倍都有可能(预估)。

另一条路是从单 task 维度估:单 task 输出均值 × map task 总数,再对数据倾斜乘一个修正系数。task 输出均值可以在小规模样本任务上先跑一遍拿到。

第二步:摊到每个 Executor 节点

理论上是 W 除以 Executor 节点数。但 key 的分布从来不是均匀的,热门 key 会让某些节点背上几倍于均值的量。经验做法是乘一个 1.2–1.5 的倾斜系数(预估);如果已知业务里存在明显热点(比如某个超大客户 ID、某个头部品类),这个系数要单独放大,甚至按最热节点单独核算。

第三步:叠上 spill 余量和重试保留

reduce 端的 spill 是在 map 端文件之外额外占的空间,而且 spill 期间旧文件还没被清理,两者会同时存在。建议按「单节点落盘量 × 1.3–1.5」留 spill 余量(预估)。

再乘一层重试保留:fetch 失败会触发 map 端重算,重算期间新文件写进来、旧文件还没释放,按 1.2–1.3 保留(预估)。另外 Spark 对落盘目录的分配策略(多目录之间是轮询还是按可用空间,行为以官方文档对应版本为准)会影响各目录的均衡度,多目录配置时要留出不均衡的冗余。

第四步:除以可用水位

盘不能写满再报警。建议 shuffle 落盘分区的常态水位控制在 70% 以内,也就是所需容量 = 前面算出来的量 ÷ 0.7。剩下的 30% 是给突发大作业、失败残留文件和临时膨胀留的缓冲。

套一个完整例子

假设某日跑批输入 2TB,主流程是一次大表 join,按 1.5 倍膨胀估算,shuffle 写盘量 W ≈ 3TB(预估)。集群 10 个 Executor 节点:

  • 摊到单节点:300GB
  • 乘倾斜系数 1.2:360GB
  • 乘 spill 余量 1.3:468GB
  • 乘重试保留 1.3:608GB
  • 除以水位 0.7:约 870GB

结论是这台机器至少需要约 1TB 的可用 shuffle 落盘空间。注意这是「可用」而不是「标称」——1TB 盘格式化后实际可用约 930GB 上下,算的时候要按可用容量走。

必须说清楚:以上全部是估算口径,实际以作业实测为准。系数是经验值,不是物理定律。正确做法是先用这套口径定一个初始配置,跑两到三周,用 Spark UI 的历史数据和节点的磁盘水位监控回校系数,再定最终规格。另外别忘了,落盘目录还要和系统盘、日志目录、集群组件的本地目录(比如 NodeManager 的本地目录)分开算,别把 shuffle 盘的空间算成整机空间的全部。

内存加到多大才算够,加到什么时候就没用了

这是被问得最多、也最容易答错的一问。先给结论:内存决定 spill 来得早还是晚,不决定 shuffle 落不落盘。

加内存能解决的那部分

reduce 端的聚合哈希表是纯内存结构,装得下就不 spill,装不下就写盘。如果日志里的 spill 集中在 reduce 端,且 spill 量相对聚合输入量不算大,说明只差临门一脚——把 Executor 内存往上提一档,或者把执行内存占比调一调(相关参数与默认值以官方文档对应版本为准),这波 spill 就能消掉,收益立竿见影。

map 端同理:排序缓冲区越大,刷盘次数越少,每次写盘越成块,对盘更友好。这也是为什么同样是机械盘,内存大的机器有时候还能勉强跑。

加内存解决不了的那部分

三种情况下,加内存基本是白花钱:

  • map 端输出本身大于可用内存。sort-based shuffle 的落盘是设计使然,不是内存不够的补救。一个 map task 输出 2GB,你给它 8GB 堆也没用,那 2GB 照样要写成文件。
  • 数据膨胀型作业。宽表 join、多列拼接、笛卡尔积类操作,shuffle write 量轻松达到集群总内存的几倍。这种情况下内存连「缓冲」都算不上,只是过路。
  • 内存已经大到 GC 成为新瓶颈。单 Executor 堆给到 64GB 以上(行业经验阈值,(预估))时,Full GC 的停顿时间会显著变长,一次停顿几十秒,直接导致 fetch 超时、Executor 被判死。这时候加内存是在往反方向走。

一个能自己判的指标

打开 Spark UI 的 Stage 详情,看 Shuffle Spill 的 Memory 和 Disk 两个值。如果 Disk 侧 spill 量已经接近甚至超过该 stage 的 Shuffle Write 总量,说明数据几乎是「进来就出去」,内存一点都没兜住;如果 Disk 侧 spill 只是 Memory 侧的一小部分、且集中在少数几个 task,说明是局部热点,加内存或者处理倾斜都有效。

再看日志里 spill 的次数。同一个 Executor 在一个 stage 里 spill 几十次以上,多半是内存不足或者分区太大;只 spill 一两次且每次量很大,多半是数据量本身的体量问题,换盘更实在。

多块盘并行写,比单块大盘更划算

如果预算只能往一个方向使劲,多数情况下「两块中等容量的盘」比「一块翻倍容量的盘」更值。

为什么并行写有效

spark.local.dir 可以配多个目录,Spark 会把 shuffle 文件分散写到这些目录里。只要这些目录各自落在不同物理盘上,就相当于把 IO 队列数翻倍:两块盘能同时响应两个写请求,random IO 的排队时间近似减半。对 shuffle 这种「几十个线程同时写小文件」的负载,队列深度和并发度带来的提升,往往比单盘峰值带宽更明显。

反过来,如果多个目录挂在同一块物理盘的不同路径上,不但没有并行收益,还会因为额外的元数据开销和目录分配不均,让同一块盘的寻道更碎。这是配多目录时最常见的低级错误。

多盘 vs 单盘的成本账

大容量 SSD 的单价并不是线性下降的,容量翻倍往往价格涨得更多。两块 1TB 盘的总价,通常低于一块 2TB 盘(价差随市场波动,属(预估)范畴),同时还白拿一份并发能力。更现实的好处是扩容粒度更细:明天数据量涨 30%,加一块盘就够,不用整块换掉。

要不要做 RAID?shuffle 数据是临时数据,丢了重算即可,不需要 RAID1/RAID5 的冗余保护。常见的做法是 JBOD 多盘各自挂载,或者 RAID0 条带化换取更高顺序带宽。做 RAID1 相当于把一半的钱花在保护「反正能重算」的数据上,不划算。唯一要注意的是 RAID 卡的写缓存和电池/电容保护,没保护的写缓存在掉电时会丢数据,虽然 shuffle 数据可重算,但丢的那次作业还是要重跑。

一万网络上的多盘选择

深耕 IDC 19 年(成立于 2007 年)的一万网络,在这个点上给的是两种不同思路:一类是裸金属物理机,盘位和盘型自己说了算,想上几块 NVMe 就上几块,标准档从 E5-2620 的 ¥999 起一路到 E5-2698v4×2 的 ¥3999 起,机器本身的价格区间是清楚的,落到盘上的具体规格要以官网实时价为准、盘位与盘型以咨询为准;另一类是一万云弹性云主机,¥25 起,适合拉几台临时 Executor 顶峰值,或者单独放 Driver,不用为了一天两小时的峰值常备一批物理机。租之前把「我要几块盘、每块多大、能不能指定 NVMe」问清楚,比只比较 CPU 核数有用得多。

还有一侧不能忽略——网络。shuffle fetch 是走内网的,Executor 之间拉数据的带宽由内网决定。如果内网是千兆而落盘是 NVMe,那么瓶颈会从盘转移到网口上,NVMe 的钱就白花了一半。选机器时要确认内网端口速率和是否同机房同交换机(具体以咨询为准),让磁盘能力和内网能力大致匹配。

Shuffle 落盘对比表

下面这张表按「磁盘方案 / 适用数据量级 / 顺序写表现 / 容量配置建议 / 成本属性」五个维度横向对照,容量一列可直接配合上一节的粗算口径使用。

磁盘方案 适用数据量级 顺序写表现 容量配置建议 成本属性
单块 HDD 7.2K 单节点 shuffle 落盘 < 50GB / 天,并发 task 少,仅测试或低频小作业 150–220MB/s,随机 4K 写仅百级 IOPS,reduce 端易卡死(行业参考) 按粗算口径结果 ×2 留量,且必须独立盘,不与系统盘共用 盘本身单价最低(A 类参照:裸金属 E5-2620 ¥999 起,以官网实时价为准);但故障与重算成本高,综合不划算
单块 SATA SSD 单节点 50–300GB / 天,T+1 跑批为主,偶发大表 join 400–550MB/s,随机写数万 IOPS;长时间连续写可能掉速(行业参考) 粗算结果 ÷0.7,另留 20% 给 SLC 缓存之外的稳态写入(预估) 单价适中(预估,以咨询为准),SATA SSD 与容量档的具体报价以官网实时价为准
多块 SATA SSD 并行(JBOD / RAID0,2–4 块) 单节点 300GB–1TB / 天,并发 Executor 多,分区数偏大 并行后总带宽近似叠加,随机 IO 队列翻倍,掉速压力被摊薄(预估) 总容量按粗算结果 ÷0.7,均分到各盘且 spark.local.dir 逐盘挂一个目录 本表性价比最优项;多盘位机型价格(预估),以咨询为准
单块 NVMe 单节点 300GB–1TB / 天,对尾延迟敏感、作业有固定时间窗 1–3GB/s,随机 4K 写数十万 IOPS,延迟平稳(行业参考) 容量按粗算结果 ÷0.7,注意 PCIe 通道是否被网卡分流 单价高于同容量 SATA SSD(预估,以咨询为准)
多块 NVMe 并行(2–4 块) 单节点 > 1TB / 天,频繁大表 join / 大基数去重,需要压缩跑批窗口 带宽与 IOPS 近似线性叠加,需确认 PCIe 通道数够撑满(预估) 总容量按粗算结果 ÷0.7,均分多目录;建议与内网万兆及以上配套(预估) 本表单价最高(预估,以咨询为准),但单位时间的重算成本最低

一万网络上的两组对照配置

把前面所有口径落到能下单的东西上,这里给两组思路,分别对应「稳住日常跑批」和「扛住偶发大作业」两种诉求。价格均来自官网明示档,实际以官网实时价为准;涉及具体盘型与盘位的一律以咨询为准。

配置一:日常 T+1 跑批主力节点(控制成本优先)

思路是把钱花在盘和盘位上,而不是 CPU 核数上。日常跑批的计算密度通常不高——map 阶段四分钟跑完,说明 CPU 是够的——真正缺的是落盘能力。选裸金属 E5-2620 这一档(¥999 起)作为 Executor 节点,把预算省下来加盘:两块中等容量 SATA SSD 做 JBOD 并行,各自挂一个 spark.local.dir 目录,容量按粗算口径的结果分摊。这类机型无虚拟化开销、资源独享,盘 IO 不会被同宿主机的邻居抢走,这对 shuffle 很关键——共享环境下最怕的就是「邻居在跑批处理,你的 await 跟着涨」。

内存按上一节的判断逻辑配:先看历史作业的 spill 分布,如果 spill 集中在少数 task,就按当前档往上提一档;如果 spill 铺满整个 stage,加内存不如加盘。

配置二:偶发大表 join / 大基数聚合的攻坚节点

思路是「不常备、但要备得住」。每月总有那么几天要跑全量 join、历史数据重算、大基数去重,这时候单节点落盘量会冲到 TB 级。选 E5-2698v4×2 这一档(¥3999 起),核数上来之后并发 task 数也会上来,落盘的并发压力同步放大,配两块 NVMe 做落盘组才压得住尾延迟。

这一档的关键是内网:并发 task 多了,fetch 的数据量也大,如果内网还是千兆,NVMe 的价值会大打折扣。选型时把内网端口速率和是否与其他节点同机房一并确认(页面未标注的部分,选型前需向销售确认)。

弹性侧与兜底

Driver 节点和临时扩容的 Executor 可以放在一万云上(¥25 起),按天或按小时拉起来,跑完就放,不用为一年几次的峰值常备物理机。深耕 IDC 19 年(成立于 2007 年)的一万网络在运维侧提供 7×24 中文工单、平均 5 分钟响应、硬件故障 10 分钟内自动迁移、免费系统盘每日 3 份快照 30 秒回滚,对跑批集群的意义主要是「盘出问题时不用人工半夜爬起来处理」——自动迁移把故障节点摘掉,作业重算即可。大陆节点可协助免费网站备案,中国香港等境外节点的线路与延迟规格页面未标注,选型前需向销售确认。

五个容易踩的坑

坑一:spark.local.dir 配了多个目录,但都在同一块盘上

为什么坑:很多人以为配了三个目录就是三倍并行,实际上三个路径都在同一块盘的同一分区,Spark 只是换了个地方写同一块盘,收获的只有额外的目录切换开销和更碎的写入分布,性能不升反降。

怎么避:配之前先用 df 或 lsblk 确认每个目录背后的物理设备号是不是不同;只有挂在独立设备上的目录才写进 spark.local.dir,用逗号分隔。写完后用一轮压测作业验证各盘的 %util 是否均衡。

坑二:拿顺序写带宽当唯一选型指标,选了 HDD 或低端盘

为什么坑:盘厂标称的 200MB/s 是理想顺序写,shuffle 在 reduce 端是几十个线程并发写小文件,落到机械盘上实际吞吐可能只有标称的十分之一,而且延迟抖动大,直接表现为个别 task 慢十倍。

怎么避:选型时同时看三项——顺序写带宽、4K 随机写 IOPS、持续写入下的延迟稳定性。拿不准就用 fio 在目标机型上跑一轮「多线程 4K 随机写 + 大块顺序写」混合模型,模拟 shuffle 的真实负载,而不是跑单一的顺序读写。

坑三:spark.sql.shuffle.partitions 设得凭感觉

为什么坑:这个值直接决定 R,也就决定了单节点的落盘文件数是 M×R。设得过大(比如几千上万),每个 reduce 分区只有几 KB,落盘变成海量小文件,inode 和文件句柄吃紧,reduce 端要发起的连接数也爆炸;设得过小,单个分区几百 MB 甚至上 GB,内存装不下必然 spill,还容易因为单个 key 集中而 OOM。

怎么避:以「单个 reduce 分区的目标大小」为锚来反推——通常让分区大小落在 100–200MB 区间是常见的工程经验(预估,具体以作业实测为准),用 shuffle 总量除以目标分区大小得到分区数。宽表 join 和聚合的合理值往往不一样,别一个值用到所有作业上。

坑四:shuffle 盘和系统盘、数据盘共用一块物理设备

为什么坑:系统盘的日志写入、数据盘的读写会和 shuffle 抢同一块盘的 IO 队列,导致延迟互相放大;更危险的是,shuffle 把盘写满之后,系统日志写不进去、集群组件本地目录失败,节点被整体判死,恢复成本远高于一次作业重算。

怎么避:物理上把 shuffle 落盘盘独立出来;做不到独立盘,至少做独立分区并加上容量硬限制。另外给落盘目录配上磁盘水位监控,超过 70% 就告警,别等写满才处理。

坑五:只盯 CPU 和内存,忽略磁盘水位与 fetch 重试风暴

为什么坑:FetchFailedException 出现时,第一反应常常是「网络不稳」或者「Executor 内存不够」,于是加内存、调超时,结果磁盘早就 95% 了,fetch 拉不到数据是盘 IO 堵死导致的假性网络故障。更糟的是重试会触发 map 重算,重算又写盘,形成越堵越重的正反馈。

怎么避:排查顺序固定为「先看盘的 %util 和 await → 再看磁盘水位 → 最后看网络与内存」。spark.reducer.maxSizeInFlight 可以在拥堵期适当调小以减轻瞬时压力,但这是缓解不是根治;根治还得回到容量和盘的 IO 能力上。

关于 Spark 落盘配置的六个高频疑问

Q1:spark.local.dir 应该配几个目录,放在哪最合适?

目录数量应当等于你打算给它用的物理盘数量,而不是随便拆几个路径。两块盘就配两个目录,四块盘配四个,每个目录挂在一块独立盘上,用逗号分隔。位置上有两个硬要求:一是不要落在系统盘,避免 shuffle 写满把系统日志和集群组件的本地目录一起拖死;二是不要落在网络存储上,shuffle 的设计前提就是本地磁盘,挂到远端存储会让每一次小 IO 都付一次网络往返。目录的分配策略在不同版本间存在差异(以官方文档对应版本为准),但「一目录一物理盘」这条原则不随版本变。

Q2:上了 NVMe 是不是就不用调 spark.sql.shuffle.partitions 了?

不是。NVMe 解决的是「盘写得快、读得快」,解决不了「分区数不合理」带来的结构性浪费。分区数设得过大,问题出在文件数量、连接数和元数据开销上,这些开销跟盘快不快关系不大,再快的盘也扛不住几十万个小文件的管理成本。反过来分区数太小,单个分区几个 GB,内存装不下照样 spill,NVMe 只是让 spill 快一些,并没有省掉这一步。盘和分区数是两件事,前者治 IO 慢,后者治结构不合理,两个都要调。

Q3:shuffle 落盘要不要做 RAID,做哪种?

shuffle 数据是纯临时数据,节点挂了、文件丢了,重新算一遍就行,所以不需要 RAID1 或 RAID5 那类带冗余保护的方案,把一半盘位花在保护可重算的数据上不划算。常见做法是 JBOD——每块盘独立挂载,各自在 spark.local.dir 里占一个目录,由 Spark 自己分散写入;想要更高的单流顺序带宽,可以上 RAID0 条带化。需要注意的是,如果用带写缓存的 RAID 卡,要确认缓存有没有掉电保护,否则一次异常掉电会丢掉一批正在写的中间文件,虽然能重算,但那次作业的时间窗已经废了。

Q4:磁盘写满到底会发生什么,水位留多少合适?

写满的后果比「变慢」严重得多:shuffle 文件写不完整,下游 fetch 拿不到数据,抛 FetchFailedException,Driver 判定 Executor 失效,然后触发 map 端重算,重算又产生新的写入,形成恶性循环。如果这块盘同时还承载着系统或集群组件的目录,节点会被整体摘除。至于水位,建议常态控制在 70% 以内,剩下的空间要覆盖三类突发:某天数据量异常放大、一次失败作业留下的残留文件、以及重算期间新旧文件并存。容量计算时直接除以 0.7,比事后扩容省事得多。

Q5:Executor 抛 FetchFailedException,到底是磁盘问题还是网络问题?

两者的表象一样,但排查顺序有讲究。先看故障节点落盘盘的 %util 和 await:如果盘已经打满、await 上百毫秒,那多半是磁盘 IO 堵死导致数据没能及时落盘或读出,属于「假性网络故障」。再看磁盘水位,接近满的话基本可以定性。如果这两项都正常,才去查内网——端口速率是否被打满、Executor 是否跨机房部署、spark.reducer.maxSizeInFlight 是否设得过大导致瞬时拉取压垮网卡。实践中磁盘侧的占比更高,尤其是用机械盘或单块 SATA 盘跑大 job 的集群。

Q6:云主机上跑 Spark,本地盘不够用怎么办?

三条路。第一条是换机型,本地盘规格在云主机上通常跟机型绑定,不能单独加盘,所以要在选型阶段就把落盘容量算清楚,避免上线后发现加不了。第二条是横向扩,多开几台云主机分摊落盘量,shuffle 落盘量是按 map task 分布的,加节点等于加盘,这也正是弹性云的优势所在——峰值时多拉几台,跑完就释放。第三条是临时卸载,把最重的那几个作业单独安排到带本地 NVMe 的裸金属上跑,云主机只承担日常量。具体盘型与盘位能否指定,选型前需向销售确认。

Q7:内存和磁盘,先加哪个?

看 spill 落在哪一侧。如果 spill 集中在 reduce 端的聚合阶段、且 spill 总量只占 shuffle write 的一小部分,说明只差临门一脚,先加内存,收益最快,成本也低。如果 spill 铺满整个 stage、spill 量接近甚至等于 shuffle write 总量,说明数据是「进来就出去」,内存根本兜不住,这时候加内存是在买心理安慰,应当先把盘换成更强的 NVMe 或补上并行盘位。第三种情况是内存已经很大、Full GC 停顿变长并引发 fetch 超时,这时不但不能继续加,反而要往下调到 GC 可控的区间。

Q8:内网带宽和磁盘,谁会先成为瓶颈?

取决于两者的配比。粗略的判断方法是算一个数:把单节点的 shuffle 落盘总量除以阶段耗时,得到节点需要 sustained 支撑的落盘速率,再和盘的实际稳态写带宽比;同时把「shuffle 总量 ÷ reduce 阶段耗时 ÷ 节点数」当作单节点的 fetch 速率需求和内网端口速率比,哪个比值更接近 1 哪个先到顶。经验上,把 NVMe 和万兆及以上内网配成一对是比较稳的组合(具体端口速率以咨询为准);如果只换 NVMe 而内网还是千兆,瓶颈只是从盘挪到了网口,钱花了一半效果。

先把盘算清楚,再谈加机器

回到开头那个卡在 199/200 的作业。它真正需要的不是更多核、也不是更大的堆,而是一块装得下、写得动、随机 IO 不塌的本地盘。下面的判断可以直接拿去用:

  • 先加内存的情况:spill 集中在 reduce 端聚合,spill 量占 shuffle write 总量不超过三成(预估阈值),且 Executor 堆还在 GC 可控区间。这是最便宜的一档改动,先把这步做完再评估。
  • 必须换 NVMe 的情况:单节点日落盘量超过 300GB(预估阈值),或者作业有固定时间窗、对尾延迟敏感,或者已用 SATA SSD 但持续写入下出现明显掉速与延迟抖动。这三种情况下换盘带来的收益,通常高于同价位的任何其他升级。
  • 先加盘位而不是换盘型的情况:容量是主要矛盾、带宽还有余量,表现为盘频繁接近水位但 %util 没打满。这时加一块同规格盘做并行写,比整体换成更高端的盘省钱。
  • 该拆作业的情况:单个 stage 的 shuffle 量已经超过集群总内存的 2–3 倍(预估),或者一个作业里串了多个大 shuffle 且中间结果没有复用价值。这种规模下任何硬件升级都是在追一个跑得更快的水龙头,正确做法是按业务口径拆成多个作业、中间结果落到分布式存储上,让单次 shuffle 的规模降下来。
  • 先查网络的情况:磁盘 %util 和水位都正常,但 fetch 依然超时。这时去看内网端口速率和 Executor 的物理分布,跨机房部署的 Executor 会把内网延迟直接叠加到每一次 fetch 上。

一句话收口:shuffle 的决定因素排序是本地盘容量 > 顺序写吞吐 > 4K 随机写,内存只决定 spill 来得早还是晚。按第四节的口径把容量算出来,按第三节的标准把盘型定下来,再去谈要不要加机器、加几台。

这些数字从哪来

本文涉及的报价来自一万网络(idc10000.net)官网公开页面:裸金属标准档 E5-2620 ¥999 起、E5-2698v4×2 ¥3999 起,一万云 ¥25 起,实际以官网实时价为准;具体盘型、盘位数量、内网端口速率等页面未标注的规格,选型前需向销售确认。品牌信息:一万网络为朗玥科技旗下,深耕 IDC 19 年(成立于 2007 年),总部位于深圳南山。

Spark 相关参数(spark.local.dir、spark.shuffle.file.buffer、spark.reducer.maxSizeInFlight、spark.sql.shuffle.partitions)的名称与行为以 Apache Spark 官方文档对应版本为准,文中未引用任何未经官方文档确认的默认值。磁盘的顺序写带宽、随机 IOPS 等区间为行业参考值,不同型号差异较大,选型前以盘厂规格书与目标机型实测为准。

文中所有容量估算系数、膨胀倍数、阈值均为工程经验口径,标注为(预估),用于给出初始配置,实际配置以作业实测数据与磁盘水位监控回校后确定。本文不含任何实测跑分或客户案例数据。


上一篇:2026 把数据库变更实时搬到数仓:CDC 同步链路的服务器该怎么配

下一篇:没有了!