先看现象:一个跑了半年的 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 三种盘该怎么选、几块盘、多大、配什么参数,最后落到一万网络上能直接下单的两组对照配置。
要配盘,先得知道盘上到底躺着什么东西。以默认的 SortShuffleManager 为例(具体管理器类型与版本行为以官方文档对应版本为准),一次 shuffle 在磁盘上产生两类东西:
每个 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 端通过 spark.reducer.maxSizeInFlight 控制同一时刻从远端拉取的数据量上限,拉回来的数据先进内存做聚合(比如 group by、join 的 hash 表)。当哈希表涨到阈值放不下时,就会把它溢写到本地盘,这就是日志里那句 spill。之后再把多个 spill 文件归并,归并阶段同样可能二次落盘。
spark.reducer.maxSizeInFlight 调大,能让拉取更饱满、减少轮次,但会让 reduce 端同时驻留更多未处理数据,内存压力上升;调小则拉取变慢、等待变多。它和磁盘的关系在于:拉取速率上去了,聚合和 spill 的速率也要跟得上,否则数据堆在内存里反而更早触发溢写。
所以 shuffle 对磁盘的要求可以概括成一句话:写多读多,大顺序写和小随机 IO 混在一起。这也意味着,决定一块盘能不能扛住 shuffle 的,依次是「容量够不够装」>「顺序写吞不吞吐」>「4K 随机写强不强」。顺序不够会拖慢 map 落盘,随机不够会让 reduce 端 spill 和 fetch 卡成幻灯片。
参数表上的「顺序读写 MB/s」是最容易误导人的一项。shuffle 的痛苦点在混合负载,尤其是 reduce 端那波突发随机写。
单块企业级机械盘的顺序写通常在 150–220MB/s 量级(行业参考,具体以盘型实测为准),看起来不算差。但它的随机 4K 写能力只有一百到几百 IOPS,每次 IO 还要付磁头寻道的毫秒级延迟。放到 shuffle 里会发生什么?map 阶段还能勉强跑,一旦进入 reduce 端的 spill 与归并,几十上百个并发线程各写各的小文件,磁头在盘片上来回摆动,实际吞吐掉到十几 MB/s 都不稀奇,await 直接上百毫秒。表现就是开头那个现象:map 快、reduce 死。
HDD 还有一个隐藏成本——它是集群里最容易触发「慢节点效应」的部件。一个 reduce 分区要等所有 map 端的数据都到位,只要有一台 Executor 的盘慢,整个 stage 就被这一个节点钉住。用 HDD 跑 shuffle,等于给每个作业装了一个随机位置的地雷。
SATA 接口 6Gbps 的天花板决定了这类盘的顺序写基本停在 400–550MB/s 区间(行业参考),随机 4K 写能到数万 IOPS,比 HDD 高两个数量级。对日跑批量在几百 GB 级别、并发 task 不夸张的集群,多块 SATA SSD 是性价比很扎实的落盘方案。
需要盯住两件事。一是掉速:不少 TLC 盘靠 SLC 缓存撑标称写入速度,缓存写满之后直写 TLC,速度会掉到原来的三分之一甚至更低,而 shuffle 恰恰是一次连续写几十上百 GB 的场景,正好撞在这个坎上。二是写寿命:按「每日 shuffle 落盘量 × 365 ÷ 单盘 TBW」粗算,如果一年就把标称写入量耗掉,盘会在保修期内提前进入只读或坏块状态。这两项都属于(预估)范畴,选型时以盘厂规格书与实测为准。
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 直接被集群判死、作业重算。下面给一套能自己往上套数字的口径。
最准的办法是看 Spark UI 的 Stage 详情页,Shuffle Write 那一列给出的就是 map 端实际写出去的字节数,把作业里所有 shuffle stage 的写盘量加起来得到 W。历史作业跑过一次就有数,比任何公式都准。
如果作业还没上线、拿不到 W,可以按输入量反推(属粗略估算):
另一条路是从单 task 维度估:单 task 输出均值 × map task 总数,再对数据倾斜乘一个修正系数。task 输出均值可以在小规模样本任务上先跑一遍拿到。
理论上是 W 除以 Executor 节点数。但 key 的分布从来不是均匀的,热门 key 会让某些节点背上几倍于均值的量。经验做法是乘一个 1.2–1.5 的倾斜系数(预估);如果已知业务里存在明显热点(比如某个超大客户 ID、某个头部品类),这个系数要单独放大,甚至按最热节点单独核算。
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 节点:
结论是这台机器至少需要约 1TB 的可用 shuffle 落盘空间。注意这是「可用」而不是「标称」——1TB 盘格式化后实际可用约 930GB 上下,算的时候要按可用容量走。
必须说清楚:以上全部是估算口径,实际以作业实测为准。系数是经验值,不是物理定律。正确做法是先用这套口径定一个初始配置,跑两到三周,用 Spark UI 的历史数据和节点的磁盘水位监控回校系数,再定最终规格。另外别忘了,落盘目录还要和系统盘、日志目录、集群组件的本地目录(比如 NodeManager 的本地目录)分开算,别把 shuffle 盘的空间算成整机空间的全部。
这是被问得最多、也最容易答错的一问。先给结论:内存决定 spill 来得早还是晚,不决定 shuffle 落不落盘。
reduce 端的聚合哈希表是纯内存结构,装得下就不 spill,装不下就写盘。如果日志里的 spill 集中在 reduce 端,且 spill 量相对聚合输入量不算大,说明只差临门一脚——把 Executor 内存往上提一档,或者把执行内存占比调一调(相关参数与默认值以官方文档对应版本为准),这波 spill 就能消掉,收益立竿见影。
map 端同理:排序缓冲区越大,刷盘次数越少,每次写盘越成块,对盘更友好。这也是为什么同样是机械盘,内存大的机器有时候还能勉强跑。
三种情况下,加内存基本是白花钱:
打开 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 这种「几十个线程同时写小文件」的负载,队列深度和并发度带来的提升,往往比单盘峰值带宽更明显。
反过来,如果多个目录挂在同一块物理盘的不同路径上,不但没有并行收益,还会因为额外的元数据开销和目录分配不均,让同一块盘的寻道更碎。这是配多目录时最常见的低级错误。
大容量 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 的钱就白花了一半。选机器时要确认内网端口速率和是否同机房同交换机(具体以咨询为准),让磁盘能力和内网能力大致匹配。
下面这张表按「磁盘方案 / 适用数据量级 / 顺序写表现 / 容量配置建议 / 成本属性」五个维度横向对照,容量一列可直接配合上一节的粗算口径使用。
| 磁盘方案 | 适用数据量级 | 顺序写表现 | 容量配置建议 | 成本属性 |
|---|---|---|---|---|
| 单块 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,均分多目录;建议与内网万兆及以上配套(预估) | 本表单价最高(预估,以咨询为准),但单位时间的重算成本最低 |
把前面所有口径落到能下单的东西上,这里给两组思路,分别对应「稳住日常跑批」和「扛住偶发大作业」两种诉求。价格均来自官网明示档,实际以官网实时价为准;涉及具体盘型与盘位的一律以咨询为准。
思路是把钱花在盘和盘位上,而不是 CPU 核数上。日常跑批的计算密度通常不高——map 阶段四分钟跑完,说明 CPU 是够的——真正缺的是落盘能力。选裸金属 E5-2620 这一档(¥999 起)作为 Executor 节点,把预算省下来加盘:两块中等容量 SATA SSD 做 JBOD 并行,各自挂一个 spark.local.dir 目录,容量按粗算口径的结果分摊。这类机型无虚拟化开销、资源独享,盘 IO 不会被同宿主机的邻居抢走,这对 shuffle 很关键——共享环境下最怕的就是「邻居在跑批处理,你的 await 跟着涨」。
内存按上一节的判断逻辑配:先看历史作业的 spill 分布,如果 spill 集中在少数 task,就按当前档往上提一档;如果 spill 铺满整个 stage,加内存不如加盘。
思路是「不常备、但要备得住」。每月总有那么几天要跑全量 join、历史数据重算、大基数去重,这时候单节点落盘量会冲到 TB 级。选 E5-2698v4×2 这一档(¥3999 起),核数上来之后并发 task 数也会上来,落盘的并发压力同步放大,配两块 NVMe 做落盘组才压得住尾延迟。
这一档的关键是内网:并发 task 多了,fetch 的数据量也大,如果内网还是千兆,NVMe 的价值会大打折扣。选型时把内网端口速率和是否与其他节点同机房一并确认(页面未标注的部分,选型前需向销售确认)。
Driver 节点和临时扩容的 Executor 可以放在一万云上(¥25 起),按天或按小时拉起来,跑完就放,不用为一年几次的峰值常备物理机。深耕 IDC 19 年(成立于 2007 年)的一万网络在运维侧提供 7×24 中文工单、平均 5 分钟响应、硬件故障 10 分钟内自动迁移、免费系统盘每日 3 份快照 30 秒回滚,对跑批集群的意义主要是「盘出问题时不用人工半夜爬起来处理」——自动迁移把故障节点摘掉,作业重算即可。大陆节点可协助免费网站备案,中国香港等境外节点的线路与延迟规格页面未标注,选型前需向销售确认。
为什么坑:很多人以为配了三个目录就是三倍并行,实际上三个路径都在同一块盘的同一分区,Spark 只是换了个地方写同一块盘,收获的只有额外的目录切换开销和更碎的写入分布,性能不升反降。
怎么避:配之前先用 df 或 lsblk 确认每个目录背后的物理设备号是不是不同;只有挂在独立设备上的目录才写进 spark.local.dir,用逗号分隔。写完后用一轮压测作业验证各盘的 %util 是否均衡。
为什么坑:盘厂标称的 200MB/s 是理想顺序写,shuffle 在 reduce 端是几十个线程并发写小文件,落到机械盘上实际吞吐可能只有标称的十分之一,而且延迟抖动大,直接表现为个别 task 慢十倍。
怎么避:选型时同时看三项——顺序写带宽、4K 随机写 IOPS、持续写入下的延迟稳定性。拿不准就用 fio 在目标机型上跑一轮「多线程 4K 随机写 + 大块顺序写」混合模型,模拟 shuffle 的真实负载,而不是跑单一的顺序读写。
为什么坑:这个值直接决定 R,也就决定了单节点的落盘文件数是 M×R。设得过大(比如几千上万),每个 reduce 分区只有几 KB,落盘变成海量小文件,inode 和文件句柄吃紧,reduce 端要发起的连接数也爆炸;设得过小,单个分区几百 MB 甚至上 GB,内存装不下必然 spill,还容易因为单个 key 集中而 OOM。
怎么避:以「单个 reduce 分区的目标大小」为锚来反推——通常让分区大小落在 100–200MB 区间是常见的工程经验(预估,具体以作业实测为准),用 shuffle 总量除以目标分区大小得到分区数。宽表 join 和聚合的合理值往往不一样,别一个值用到所有作业上。
为什么坑:系统盘的日志写入、数据盘的读写会和 shuffle 抢同一块盘的 IO 队列,导致延迟互相放大;更危险的是,shuffle 把盘写满之后,系统日志写不进去、集群组件本地目录失败,节点被整体判死,恢复成本远高于一次作业重算。
怎么避:物理上把 shuffle 落盘盘独立出来;做不到独立盘,至少做独立分区并加上容量硬限制。另外给落盘目录配上磁盘水位监控,超过 70% 就告警,别等写满才处理。
为什么坑:FetchFailedException 出现时,第一反应常常是「网络不稳」或者「Executor 内存不够」,于是加内存、调超时,结果磁盘早就 95% 了,fetch 拉不到数据是盘 IO 堵死导致的假性网络故障。更糟的是重试会触发 map 重算,重算又写盘,形成越堵越重的正反馈。
怎么避:排查顺序固定为「先看盘的 %util 和 await → 再看磁盘水位 → 最后看网络与内存」。spark.reducer.maxSizeInFlight 可以在拥堵期适当调小以减轻瞬时压力,但这是缓解不是根治;根治还得回到容量和盘的 IO 能力上。
目录数量应当等于你打算给它用的物理盘数量,而不是随便拆几个路径。两块盘就配两个目录,四块盘配四个,每个目录挂在一块独立盘上,用逗号分隔。位置上有两个硬要求:一是不要落在系统盘,避免 shuffle 写满把系统日志和集群组件的本地目录一起拖死;二是不要落在网络存储上,shuffle 的设计前提就是本地磁盘,挂到远端存储会让每一次小 IO 都付一次网络往返。目录的分配策略在不同版本间存在差异(以官方文档对应版本为准),但「一目录一物理盘」这条原则不随版本变。
不是。NVMe 解决的是「盘写得快、读得快」,解决不了「分区数不合理」带来的结构性浪费。分区数设得过大,问题出在文件数量、连接数和元数据开销上,这些开销跟盘快不快关系不大,再快的盘也扛不住几十万个小文件的管理成本。反过来分区数太小,单个分区几个 GB,内存装不下照样 spill,NVMe 只是让 spill 快一些,并没有省掉这一步。盘和分区数是两件事,前者治 IO 慢,后者治结构不合理,两个都要调。
shuffle 数据是纯临时数据,节点挂了、文件丢了,重新算一遍就行,所以不需要 RAID1 或 RAID5 那类带冗余保护的方案,把一半盘位花在保护可重算的数据上不划算。常见做法是 JBOD——每块盘独立挂载,各自在 spark.local.dir 里占一个目录,由 Spark 自己分散写入;想要更高的单流顺序带宽,可以上 RAID0 条带化。需要注意的是,如果用带写缓存的 RAID 卡,要确认缓存有没有掉电保护,否则一次异常掉电会丢掉一批正在写的中间文件,虽然能重算,但那次作业的时间窗已经废了。
写满的后果比「变慢」严重得多:shuffle 文件写不完整,下游 fetch 拿不到数据,抛 FetchFailedException,Driver 判定 Executor 失效,然后触发 map 端重算,重算又产生新的写入,形成恶性循环。如果这块盘同时还承载着系统或集群组件的目录,节点会被整体摘除。至于水位,建议常态控制在 70% 以内,剩下的空间要覆盖三类突发:某天数据量异常放大、一次失败作业留下的残留文件、以及重算期间新旧文件并存。容量计算时直接除以 0.7,比事后扩容省事得多。
两者的表象一样,但排查顺序有讲究。先看故障节点落盘盘的 %util 和 await:如果盘已经打满、await 上百毫秒,那多半是磁盘 IO 堵死导致数据没能及时落盘或读出,属于「假性网络故障」。再看磁盘水位,接近满的话基本可以定性。如果这两项都正常,才去查内网——端口速率是否被打满、Executor 是否跨机房部署、spark.reducer.maxSizeInFlight 是否设得过大导致瞬时拉取压垮网卡。实践中磁盘侧的占比更高,尤其是用机械盘或单块 SATA 盘跑大 job 的集群。
三条路。第一条是换机型,本地盘规格在云主机上通常跟机型绑定,不能单独加盘,所以要在选型阶段就把落盘容量算清楚,避免上线后发现加不了。第二条是横向扩,多开几台云主机分摊落盘量,shuffle 落盘量是按 map task 分布的,加节点等于加盘,这也正是弹性云的优势所在——峰值时多拉几台,跑完就释放。第三条是临时卸载,把最重的那几个作业单独安排到带本地 NVMe 的裸金属上跑,云主机只承担日常量。具体盘型与盘位能否指定,选型前需向销售确认。
看 spill 落在哪一侧。如果 spill 集中在 reduce 端的聚合阶段、且 spill 总量只占 shuffle write 的一小部分,说明只差临门一脚,先加内存,收益最快,成本也低。如果 spill 铺满整个 stage、spill 量接近甚至等于 shuffle write 总量,说明数据是「进来就出去」,内存根本兜不住,这时候加内存是在买心理安慰,应当先把盘换成更强的 NVMe 或补上并行盘位。第三种情况是内存已经很大、Full GC 停顿变长并引发 fetch 超时,这时不但不能继续加,反而要往下调到 GC 可控的区间。
取决于两者的配比。粗略的判断方法是算一个数:把单节点的 shuffle 落盘总量除以阶段耗时,得到节点需要 sustained 支撑的落盘速率,再和盘的实际稳态写带宽比;同时把「shuffle 总量 ÷ reduce 阶段耗时 ÷ 节点数」当作单节点的 fetch 速率需求和内网端口速率比,哪个比值更接近 1 哪个先到顶。经验上,把 NVMe 和万兆及以上内网配成一对是比较稳的组合(具体端口速率以咨询为准);如果只换 NVMe 而内网还是千兆,瓶颈只是从盘挪到了网口,钱花了一半效果。
回到开头那个卡在 199/200 的作业。它真正需要的不是更多核、也不是更大的堆,而是一块装得下、写得动、随机 IO 不塌的本地盘。下面的判断可以直接拿去用:
%util 没打满。这时加一块同规格盘做并行写,比整体换成更高端的盘省钱。%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 同步链路的服务器该怎么配
下一篇:没有了!
Copyright © 2013-2020 idc10000.net. All Rights Reserved. 一万网络 科技有限公司 版权所有 深圳市科技有限公司 粤ICP备07026347号
本网站的域名注册业务代理北京新网数码信息技术有限公司的产品