Apache Uniffle:根治Spark Shuffle磁盘瓶颈与任务重算的远程Shuffle引擎 凌晨两点值班群里一条告警一个跑了三个多小时的 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 ServiceRSS核心思路大致相同计算存储分离、独立 Shuffle 服务集群。但落到工程实现上差别很大。我整理了一张对比表对比维度本地 ShuffleFacebook RSSApache Uniffle数据存储位置计算节点本地盘独立 RSS 集群HDFS独立 Server本地盘/HDFS/对象存储支持引擎通用主要 SparkSpark、MapReduceTez 生态持续完善多副本无有限支持原生支持多副本存储可插拔不适用主要绑定 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.managerorg.apache.spark.shuffle.RssShuffleManager spark.rss.coordinator.quorumcoordinator1:19999,coordinator2:19999 spark.rss.storage.typeLOCALFILE spark.serializerorg.apache.spark.serializer.KryoSerializer提交命令里带上 client 相关 jarspark-submit \ --class com.example.MyJob \ --jars /path/to/rss-client-spark3-shaded.jar \ --conf spark.shuffle.managerorg.apache.spark.shuffle.RssShuffleManager \ --conf spark.rss.coordinator.quorumcoordinator-host:19999 \ --conf spark.rss.storage.typeLOCALFILE \ 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.capacityServer 接收缓冲区总容量按 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 WriteServer 内存被打爆。排查路径先去看 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 左右的真实作业试试水用数据说话再决定要不要铺开。