凌晨两点,值班群里一条告警:一个跑了三个多小时的 Spark 离线任务,卡在 Shuffle 阶段,磁盘写满,executor 反复重试最终还是失败。点开 Spark UI 一看,上百个 executor 所在节点的磁盘分布极不均匀,有的盘用了 60%,有的盘已经 99%。这种场景在跑日均 PB 级数据的集群里几乎是家常便饭。当时我脑子里冒出来的念头是:如果这批 Shuffle 数据不是写在计算节点的本地盘,而是写到一个独立、专门的 Shuffle 服务上,很多问题从一开始就不会发生。
这就是今天想认真聊的 Apache Uniffle。它原本是腾讯在自研过程中沉淀下来的一套远程 Shuffle 服务,后来捐献给 Apache 基金会,定位是做一个统一的远程 Shuffle 引擎。核心思路一句话:把 Shuffle 的数据存储和计算节点彻底解耦,Map 阶段产出的中间数据不落本地盘,而是写到独立的 Shuffle Server 上,Reduce 阶段再从 Server 拉取。如果你正在用 Spark 或 MapReduce 跑大规模离线任务,或者经常被 Shuffle 阶段的磁盘打满、节点宕机重算、数据倾斜这几件事折磨,这篇文章会从原理、架构、部署到调优捋一遍,读完你基本能判断自己的团队要不要引入它,以及大概怎么落地。
1. Shuffle 为什么成了分布式计算的头号瓶颈
1.1 先搞清楚 Shuffle 到底在干什么
Shuffle 翻译成中文叫"洗牌",在分布式计算框架里的含义是:上游任务产出的数据,如何按照目标分区重新组织,并交给下游任务消费。举个 Spark 里最简单的 groupByKey 例子:上游 map 阶段每个 task 处理一部分输入,产出的 (key, value) 会被哈希到不同的下游分区。问题是,上游 1000 个 task 手里各自都有一部分属于"分区 0"的数据,下游负责分区 0 的 reduce task,必须把散落在所有上游 task 里的数据全部集齐,才能开始聚合计算。
这个过程天然是全集群范围的、跨节点跨网络的数据搬运,而且搬运量通常比计算量还大几个量级。这也决定了 Shuffle 阶段的性能天花板往往不在 CPU,而在磁盘 IO、网络带宽和文件系统元数据上。很多 Spark 作业跑得慢,不是计算逻辑有多复杂,而是时间全耗在数据搬来搬去上了。
1.2 本地 Shuffle 的三宗罪
先说小文件问题。经典的本地 Shuffle 落地方式是:每个上游 task 给每个下游分区写一个文件,M 个 map task 对应 R 个 reduce 分区,就会产生 M×R 个文件。1000 个 map task、1000 个分区,就是一百万个文件。哪怕每个文件只有几十 KB,对文件系统 inode 和元数据服务都是巨大压力。Spark 的排序版 Shuffle 会做一定的合并,但文件的量级依然很吓人,作业一跑起来,整个集群到处是细碎的小文件在读写在。
第二宗罪是磁盘热点。数据分布天然不均匀,哪怕上游 task 处理的数据量差不多,经过 hash 分区之后落到某个下游分区的数据也可能明显高于平均值,这就是数据倾斜的常见形态。数据写到哪台节点、哪块盘,完全取决于 task 被调度到哪,而不是哪块盘有空间。于是经常出现:某些磁盘被打满,作业报错;旁边一堆磁盘空闲,资源利用率极度畸形。我做过的几次事故复盘里,磁盘分不均匀导致的 Shuffle 失败占比相当高。
第三宗罪最伤筋动骨:任务重算放大。本地 Shuffle 模式下,下游要的数据散落在上游 task 所在节点的本地磁盘上。一旦某台节点磁盘坏块、机器宕机、或者被混部系统强杀,上面多个 map task 的 Shuffle 中间结果全部丢失,Spark 只能重新调度这些上游 task 重算一遍。大集群里,一个节点宕机引发成百上千个 task 重算的场景,我见过太多次,这个代价极其昂贵,而且往往是雪崩式的。
1.3 集群规模一大,问题的性质就变了
小集群、小数据量的时候,上面这些问题咬咬牙都能忍。但数据量从 GB 级涨到 TB 级、PB 级,集群从几十台涨到上千台,Shuffle 的失败率会显著上升,作业耗时里 Shuffle 占比动不动就超过一半。到了这个阶段,问题已经不是"这个任务能不能跑完",而是"整个集群有多少资源被白白耗在了数据拷贝和重算上"。这也是为什么近两年各家大厂都在做远程 Shuffle——要么自研,要么直接落到 Uniffle 这样的开源方案上。远程 Shuffle 不是锦上添花,而是在大数据量场景下维持作业稳定性和集群利用率的一个刚需组件。
2. Uniffle 做了什么:把 Shuffle 从计算节点挪到独立服务
2.1 核心设计思想:计算存储分离
Uniffle 的思路说白了就是一句话:Shuffle 本质上是存储问题,不是计算问题,那就别让计算节点背这个锅。Map 阶段的输出不要写到执行节点本地盘,而是通过网络发到一组专门的 Shuffle Server 上。Shuffle Server 只干一件事:收数据、按分区存数据、等下游来取。计算节点自身变得"无状态",节点挂了不需要重算上游,因为数据根本不在它本地。
这个思路用生活化的类比很好理解:以前每个饭店各自囤菜,厨师做菜前得花大量时间去自己的后厨翻找食材;Uniffle 相当于建了一个中央仓库,所有供应商把半成品送进仓库,各个饭店做菜时直接从仓库取货,任何一家饭店的后厨失火,都不会影响整个供应链交付。数据不跟着计算节点走,是这套方案最本质的转变。
2.2 一次完整的数据流转过程
一个 Spark 作业接入 Uniffle 后,完整链路大概是这样的:
- 作业启动时,Driver 通过 ShuffleManager 向 Coordinator 注册,申请一批可用的 Shuffle Server。
- Coordinator 根据各 Server 的负载情况,给这个作业分配一份 Server 列表,负载高的节点会被跳过。
- 每个 Map task 写完输出后,按下游分区把数据切分成更小的 block,通过异步 Netty 客户端推送到分配的 Server。
- Shuffle Server 收到 block 后写入存储层,并按配置决定是否做多副本复制。
- Reduce task 启动时,向 Coordinator 获取 Server 列表,再从对应 Server 拉取属于自己的分区数据。
- 拉全一个分区的所有 block 后,Reduce 端做合并、排序、聚合,进入正式计算。
第 5 步是关键:Reduce 不再依赖上游 task 是否还存活,只要 Server 上的数据还在,它就能拉着走。上游节点的任何故障,都与下游的数据读取解耦了。Server 端的数据组织方式是按 partition 建数据文件,再用独立的 index 文件记录每个 block 的偏移量。下游拉取时先查 index,再精确定位读取对应的 block,避免了整文件扫描带来的 IO 浪费。
2.3 Coordinator、Shuffle Server、Client 各司其职
Uniffle 整个体系里就三种角色,结构非常干净。
Coordinator(协调者):相当于调度中枢,但它不碰数据本身,只负责三件事。一是维护集群里所有 Shuffle Server 的存活状态和负载信息;二是给每个作业分配 Server 列表,分配时避开高负载节点和黑名单节点;三是做作业注册、心跳和过期清理。Coordinator 可以部署多个组成 HA,借助 ZooKeeper 做选主,避免单点。
Shuffle Server(洗牌服务节点):数据面的核心。每个 Server 上挂两类线程,一类负责接收上游推送的 block 并写入存储,一类负责响应下游读取请求。接收的数据先落到内存 buffer,再异步刷到存储层。存储层通过配置项支持本地磁盘(LOCALFILE)、HDFS、对象存储(Ozone/S3)等多种形态。Server 定期向 Coordinator 上报心跳,如果一段时间失联,Coordinator 就会把它从分配列表里摘掉。
Client(客户端):就是打进 Spark / MapReduce 作业里的那部分代码,以 ShuffleManager 插件的形式运行在 Driver 和 Executor 里。它负责向 Coordinator 注册、请求 Server 列表、把 Map 输出切块推送、从 Server 拉取数据。对普通用户来说,接入动作就是替换 ShuffleManager 并配几个参数,业务代码一行都不用改,这一点对推广落地非常友好。
2.4 多副本与容错设计
Uniffle 的容错设计比我见过的不少自研方案要完整。它可以配置同一份数据在多个 Server 上写多份副本,某个 Server 挂掉之后,下游依然能从其他副本拉数据,副本数和一致性级别都可调。这相当于给 Shuffle 数据上了保险,节点故障不再意味着上游重算。
另一层保障是数据校验。传输和存储过程中会带 checksum,下游拉取时校验数据完整性,一旦发现损坏可以触发重拉或报错,而不是默默算出一个错误结果。再加上动态拉黑机制:Server 写失败或长时间 GC 停顿、心跳异常时,Coordinator 会把它拉黑,新的作业不会往它上面分配数据。这套机制保证了集群里即使有一部分节点在慢慢变烂,整体作业依然能维持可接受的完成率。
3. 同类远程 Shuffle 方案横向对比
3.1 Facebook RSS 与 Uniffle 的异同
做远程 Shuffle 的,Uniffle 不是第一家。Facebook 内部就有一套 Remote Shuffle Service(RSS),核心思路大致相同:计算存储分离、独立 Shuffle 服务集群。但落到工程实现上,差别很大。我整理了一张对比表:
| 对比维度 | 本地 Shuffle | Facebook RSS | Apache Uniffle |
|---|---|---|---|
| 数据存储位置 | 计算节点本地盘 | 独立 RSS 集群(HDFS) | 独立 Server(本地盘/HDFS/对象存储) |
| 支持引擎 | 通用 | 主要 Spark | Spark、MapReduce,Tez 生态持续完善 |
| 多副本 | 无 | 有限支持 | 原生支持多副本 |
| 存储可插拔 | 不适用 | 主要绑定 HDFS | 本地盘、HDFS、Ozone、S3 等多形态 |
| 动态分配场景 | 依赖外部 Shuffle Service | 支持 | 支持 |
| 社区形态 | 不适用 | 内部分享为主 | Apache 社区持续迭代 |
Uniffle 是从腾讯大规模生产环境反复打磨过的,它把远程 Shuffle 的能力做成通用引擎开源,对广大的开源用户来说价值更直接。因为它不只是解决"有没有远程 Shuffle"的问题,还解决了"接到自己的技术栈里是否顺手"的问题。
3.2 和 Spark 内置 Push-based Shuffle 的差异
近两年 Spark 3.2 之后推出了内置的 Push-based Shuffle,很多人拿它和 Uniffle 比较。我的看法是:Spark 内置方案是不错的基础设施,但它本质上还是在 executor 和 ESS 服务的框架内做优化,和"计算存储分离"这个目标还有距离。内置方案需要单独部署 ESS 服务,支持的功能和可调维度也更有限。Uniffle 的独立性带来的直接好处是:Shuffle 集群可以和计算集群完全分开扩容。计算节点缩容、故障、重启,都不用担心 Shuffle 数据丢失,运维边界非常清晰,这对经常弹性扩缩容的云原生环境尤其重要。
3.3 什么情况下值得引入远程 Shuffle
不是所有团队都需要远程 Shuffle,这个我得说在前面。我自己判断是否引入的标准有这么几条:
- 单作业 Shuffle 数据量在几百 GB 以上,且 Shuffle 阶段耗时占总作业耗时 40% 以上;
- 集群节点故障率偏高,混部或者裸金属场景下经常出现节点宕机,触发大规模 task 重算;
- 数据倾斜在业务层面无法根本消除,需要靠运行时机制来缓解;
- 跑在 Kubernetes 上,Pod 销毁和重建频繁,本地 Shuffle 数据很难随 Pod 生命周期保留。
如果只是几十台机器的中小集群,作业量在几百 GB 以内,先把 Spark 自身的参数调好、把数据倾斜治理掉,性价比可能更高。远程 Shuffle 要额外承担一批机器的成本和运维复杂度,这笔账得算清楚,不是越先进就越该上。
4. 部署与接入实践:用 Spark 3 跑通最小集群
4.1 环境准备
部署 Uniffle 本身不复杂,但对环境有几点要求。首先是 JDK,官方推荐 JDK 8,部分新版本可以跑在 JDK 11 上;然后是 Hadoop 客户端,哪怕你只用 LOCALFILE 存储,也建议把 Hadoop classpath 备好,很多辅助功能会依赖它。获取安装包最省事的办法,是直接用官方 Release 页面里打好的 tar.gz,也可以拉源码用项目自带的build_distribution.sh构建:
./build_distribution.sh构建前把 JDK 8 和 Maven 环境配好,输出目录下会得到完整的发行包,里面有 bin、conf、lib 整套东西。如果只是给 Spark 任务接入用,不部署服务端,那直接从 Maven Central 拿 client jar 就行。这里就体现出 Maven 生态的方便之处了:Uniffle 的客户端依赖是发到中央仓库的,你在自己的工程里引入对应版本的 rss-client 即可,版本号按你部署的 Uniffle 版本来。
4.2 部署 Coordinator 和 Shuffle Server
一个最小集群至少需要一个 Coordinator 和一个 Shuffle Server。两个角色的配置分别在conf/coordinator.conf和conf/server.conf里。Coordinator 参考配置:
rss.coordinator.rpc.server.port 19998 rss.coordinator.netty.server.port 19999 rss.jetty.http.port 19980 rss.coordinator.app.expired.withoutHeartbeat 60000Shuffle Server 参考配置(以本地盘存储为例):
rss.rpc.server.port 19988 rss.server.netty.port 19989 rss.storage.type LOCALFILE rss.server.flush.thread 32 rss.server.buffer.capacity 40g rss.server.read.buffer.capacity 20g rss.server.hadoop.cfg.dir /etc/hadoop/conf这里rss.server.buffer.capacity是接收缓冲区总容量,决定 Server 能在内存里同时承接多少上游推送数据,配小了容易触发反压,拖慢整个作业;rss.server.read.buffer.capacity对应读取侧容量,主要影响下游拉取效率。LOCALFILE 模式下,数据落在部署 Server 的机器本地磁盘,还可以通过数据目录相关配置指定多块盘,把 IO 压力摊开。不同版本的参数名可能略有差异,以你下载版本的官方文档为准。
启动方式很直接:
bin/start-coordinator.sh bin/start-shuffle-server.sh生产环境我强烈建议至少两个 Coordinator,用 ZooKeeper 做 HA。单个 Coordinator 虽然也能跑,但一旦它挂了,新作业全部注册失败,整个批处理链路就断了。这个单点是真的没必要留。
4.3 Spark 侧接入
Spark 接入是 Uniffle 做得最顺滑的地方。不需要改业务代码,提交作业的时候把 ShuffleManager 换成 Uniffle 的,再指定 Coordinator 地址即可。关键参数:
spark.shuffle.manager=org.apache.spark.shuffle.RssShuffleManager spark.rss.coordinator.quorum=coordinator1:19999,coordinator2:19999 spark.rss.storage.type=LOCALFILE spark.serializer=org.apache.spark.serializer.KryoSerializer提交命令里带上 client 相关 jar:
spark-submit \ --class com.example.MyJob \ --jars /path/to/rss-client-spark3-shaded.jar \ --conf spark.shuffle.manager=org.apache.spark.shuffle.RssShuffleManager \ --conf spark.rss.coordinator.quorum=coordinator-host:19999 \ --conf spark.rss.storage.type=LOCALFILE \ my-job.jar提示:client jar 的版本要和你用的 Spark 大版本、Uniffle 服务端版本都匹配。Spark 3.1、3.2、3.3、3.4 对应的 client 可能不是同一个 artifact,版本错配最常见的报错是 ShuffleManager 类找不到或者序列化异常。
YARN 模式下,--jars会被上传到 DistributedCache,由 NodeManager 加载到 executor classpath,一般没问题。但如果你的集群自定义了 Spark 的 classpath 覆盖逻辑,记得手工把 client jar 加进去。这个问题我接手别人的集群时踩过,现象就是作业一进 Shuffle 阶段立刻报类找不到,排查了老半天才发现是用户自己定制了 classpath。
4.4 验证接入是否生效
接入之后怎么确认真的生效了?不要只看作业跑完就放心。我会依次核对几个信号:
- 日志里出现
RssShuffleManager相关的初始化信息; - Spark UI 的 Shuffle 页面里,Shuffle Write 不再指向本地临时目录,而是能看到具体 Shuffle Server 的地址标识;
- Shuffle Server 的日志和指标上有对应的 appId 注册、block 写入量增长;
- Coordinator 的 Web 页面或接口上能看到这个作业分配到的 Server 列表。
第一次建议用一个小作业验证,跑完后把各环节日志对照一遍,再逐步放大规模。别一上来直接压到 TB 级,否则出了问题都不知道是配置问题还是负载问题,排查成本会很高。
5. 参数调优与生产环境踩坑
5.1 必调参数清单
部署是一回事,跑得好是另一回事。下面几个参数是我实测下来影响最明显的,按场景整理成表格:
| 参数 | 作用 | 我的建议 |
|---|---|---|
| rss.server.buffer.capacity | Server 接收缓冲区总容量 | 按 Server 机器内存的 40%~60% 配,太小会频繁触发反压 |
| rss.server.read.buffer.capacity | 读取缓冲区容量 | 磁盘慢、读延迟高的场景调大 |
| rss.server.flush.thread | 刷盘线程数 | 高吞吐场景至少 32 起步 |
| spark.rss.client.send.size.limit | 客户端单次推送数据大小上限 | 默认值偏保守,网络好的集群可以适当调大 |
| spark.rss.client.read.buffer.size | 客户端读取缓冲大小 | 下游拉取大分区时调大,能明显减少 RPC 次数 |
| rss.client.write.buffer.size | 客户端写缓冲大小 | 决定 map 端聚合 block 的粒度,太小会产生过多小 block |
调参的要义不是把每个参数都调到最大,而是找到当前集群资源约束下的平衡点。网络带宽充足、磁盘 IO 跟得上,就放大 buffer;磁盘是瓶颈,就去调刷盘线程和 block 大小。盲目加大内存 buffer 而不考虑 GC,反而会把 Server 拖垮。
5.2 常见故障与排查思路
这一节我把真机上踩过的几个坑列出来,每个都附带排查思路,比直接给结论有用。
坑一:作业卡在 Shuffle Write,Server 内存被打爆。排查路径:先去看 Server 上buffer.capacity是不是远小于实际写入量;再看是不是客户端并发推送太快、Server 刷盘跟不上。我的处理顺序是:先加flush.thread,再调大buffer.capacity,如果还不够,调低客户端的推送粒度,给 Server 端减压。排查顺序建议是"先看 Server 指标,再看磁盘 IO,最后才动客户端参数",很多人一上来就调大 buffer,等于把问题往后推。
坑二:Reduce 拉取数据超时。从两端看:一是 Server 端磁盘 IO 是否饱和,读请求排队严重;二是客户端读取缓冲如果配太小,同一个分区数据要分很多次 RPC 才能拉完,每次都重新建连接,时间消耗自然上去了。我遇到过一个案例,把 read.buffer.size 从 8m 调到 32m 之后,整个 stage 耗时直接降了 30%。这个优化成本几乎为零,收益却很直观。
坑三:Coordinator 把 Server 拉黑,作业大面积失败。往往是 Server 和 Coordinator 之间的心跳网络不稳定,或者 Server 因为 GC 停顿太久导致心跳超时。这时候别急着怀疑 Uniffle 本身,先看一眼 GC 日志和网络丢包率。拉黑是保护机制而不是故障本身,把网络和 JVM 参数调稳,问题自然消失。
坑四:动态分配下 executor 被回收,Shuffle 数据"看起来丢了"。早期版本配合 Spark 动态分配时容易出现,远程 Shuffle 的数据本来不在 executor 本地,正常不该丢数据,但客户端推送 block 是异步的,如果 executor 在 block 还没全部推送完成就被回收,就可能丢块。解决办法是给 executor 回收设置一个宽限期,或者升级到已经处理这类场景的新版本。生产环境稳妥起见,可以先关掉动态分配跑一个月,稳定后再逐步打开。
5.3 一点经验体会
最后说点个人感受。Uniffle 不是我见过最"炫"的组件,它解决的问题非常底层、非常痛,但没有太多花哨的设计,架构干净、思路清晰。它真正让我觉得值的地方,不只是性能提升,而是运维模型的简化:以前 Spark 计算节点挂一台,我担心的是重算风暴、磁盘写满、任务雪崩;现在计算节点出问题,影响面被 Shuffle Server 集群隔离了,我只需要关注 Server 集群的健康度就好。这个心智负担的减轻,在凌晨被告警叫醒的时候感受特别明显。
另外我会建议尽早把 Uniffle 的关键监控指标接入自家监控体系,重点是每个 Server 的接收速率、刷盘延迟、block 失败率。作业出问题时,这些指标定位故障的速度,远比你逐条翻日志快得多。如果你的集群也踩在我前面说的那几个痛点上,不妨先在测试集群搭一套最小 Uniffle,拿一个 1TB 左右的真实作业试试水,用数据说话,再决定要不要铺开。