ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Uniffle:把Spark Shuffle搬进远程数据中转站

2026/9/16 21:40:40 拓冰建站 浏览量
Uniffle:把Spark Shuffle搬进远程数据中转站 从入行第一天起我就被反复教育大数据作业的稳定性一半看 Shuffle。以前我不信直到线上有个每天跑的 Spark 批任务DataNode 磁盘被打满、上游重算、下游超时一连三天凌晨被电话叫醒。排查到最后问题全都堆在 Shuffle 写入那一环。后来我把这套作业接到了 Apache Uniffle 上才算是把这块数据搬运动脉给彻底理顺了。今天这篇就当是每天认识一个组件系列里的一篇从使用者的角度把 Uniffle 拆开讲讲它解决的是哪类问题、体系怎么设计的、接入时真正要留神的细节又是什么。1. Shuffle 为什么总是那个背锅侠1.1 一次完整 Shuffle 到底在搬什么数据不管是 Spark 还是 MapReduceShuffle 的本质都是把上游任务产生的、按照 Key 分散在多个节点上的中间结果重新划分给下游任务去消费。举个具体例子假设你有 100 个 Map 任务产出 200 个分区Partition下游有 200 个 Reduce 任务。那每个 Map 任务在输出时就要按哈希规则把数据写进对应的 200 个分区文件里每个 Reduce 任务启动后需要从 100 个 Map 任务的输出中把属于自己的那 1/200 数据全部拉回来再进行合并排序。这个各自写一堆、再互相拉一把的过程就是 Shuffle。理解它只需要抓住两个动作一个是 Map 端的 Shuffle Write一个是 Reduce 端的 Shuffle Read。框架再怎么包装底层逃不开这两个阶段。在 Spark 的传统实现里Shuffle Write 产生的数据会先落在 Executor 本地磁盘Reduce 端通过 BlockManager 或 External Shuffle Service 去拉取。这个设计本身没有问题在小规模作业上甚至很高效——数据本地读取不走网络。但问题在于当作业规模变大、Shuffle 数据量从几十 GB 涨到几个 TB 时本地落盘这套逻辑就开始露出疲态了。1.2 传统本地落盘方案的四宗原罪第一宗罪是计算与存储强耦合。Executor 既要跑用户代码又要当 Shuffle 数据的临时仓库。磁盘 IO、内存、CPU 在计算和数据中转两件事之间互相争抢一旦 Shuffle 数据激增最先被打垮的反而是那些本来要专注计算的节点。我们的线上事故就是这么来的磁盘满BlockManager 写不进去任务失败然后重试然后继续写不进去雪崩式恶化。第二宗罪是数据可靠性太依赖节点。在 MapReduce 时代Shuffle 数据是落在 Task 运行的节点上的Task 挂了数据就没了只能整个 Task 重跑。Spark 引入了 External Shuffle ServiceESS把 Shuffle 数据的服务进程从 Executor 中挪出来Executor 挂了数据还在这确实解决了一部分问题。但 ESS 依然是节点本地的守护进程磁盘损坏、机器宕机数据一样灰飞烟灭。第三宗罪是长尾效应被放大。Reduce 端要等所有上游 Map 的输出都准备好才能拉全数据。任何一个 Map 任务因为磁盘抖动变慢整个 Stage 就要陪它等。在本地落盘模式下这种木桶效应几乎无法回避。第四宗罪更直接——运维没法预估容量。你没法提前告诉 YARN 或 Kubernetes这个作业将产生 2TB Shuffle 数据只能靠经验给每台机器预留空间。预留多了浪费预留少了直接跑挂。在容器化、资源池化的环境下这个问题尤其尖锐Pod 的本地盘本来就不是为海量中间数据设计的。这四宗罪叠加在一起让把 Shuffle 数据从计算节点上拿出去成了必然趋势。Uniffle 就是在这个背景下进入我视野的。2. Uniffle 的核心思路建一个统一的数据中转站2.1 从腾讯内部 RSS 到 Apache 社区项目Uniffle 的前身是腾讯内部的远程 Shuffle 服务Remote Shuffle ServiceRSS在腾讯内部经过了大规模生产环境的检验2022 年捐赠给 Apache 基金会进入孵化器后来有了 Apache Uniffle 这个名字。它的核心思想很直接既然本地落盘毛病这么多那就单独搞一个集群专门负责接收、存储和提供 Shuffle 数据。计算节点Executor/Task只负责把数据推给这个中转站下游计算的时候再来中转站取数据。这样计算资源池和 Shuffle 存储资源池彻底分开各自弹性伸缩互不拖累。2.2 三个角色一台戏Uniffle 的整体架构里有三个必须搞清楚的组件角色对应进程职责类比Coordinator无状态服务通常部署多个管理 ShuffleServer 的注册和心跳维护集群资源视图为每个 Shuffle 分配服务器客户端读写前先来它这里问路调度台ShuffleServer有状态服务负责实际数据存取接收 Map 端推送的数据在内存缓冲中暂存按策略刷到 HDFS 或本地磁盘响应 Reduce 端的拉取请求物流仓库Client内嵌在 Spark/MR 中的插件拦截 Job 的 ShuffleManager 调用把数据写入改成推送给 ShuffleServer把读取改成从 ShuffleServer 拉取发货员和收货员Coordinator 的部署非常轻可以理解成一个活地图。它不存数据只存元数据和状态。ShuffleServer 启动时向 Coordinator 注册之后不断上报心跳当一个 Spark 应用要注册 Shuffle 时Coordinator 会根据服务器的资源使用情况给这个 Shuffle 分配一组服务器并返回给客户端。这样客户端就知道该把数据推给谁了。2.3 Map 端推送与 Reduce 端拉取的完整链路先看写入链路。在 Spark 里一旦你把 ShuffleManager 替换成 Uniffle 提供的 RssShuffleManager原来走本地 BlockManager 的写入路径就变了Map 任务输出的数据先进入客户端内存缓冲按 64KB 一个批次进行切分缓冲攒够一批就通过 gRPC/Netty 推送给 Coordinator 指定的 ShuffleServerShuffleServer 收到数据后先放入自己的内存缓冲缓冲水位达到阈值后再批量刷到存储层HDFS 或本地磁盘为了容错写入通常配了多副本多个 ShuffleServer 都确认收到后客户端才认为这次写入成功。再看读取链路。Reduce 任务启动后通过 RssShuffleManager 向 Coordinator 询问我需要的分区数据在哪些 ShuffleServer 上拿到地址列表后Reduce 端并发地从这些 ShuffleServer 拉取属于自己的分区数据拉数据时优先命中的是 ShuffleServer 内存里还热着的块内存里没有的再从存储层读数据到 Reduce 端后做常规的排序、合并进入后续计算。这个流程把原来每个 Executor 自己又是生产者又是仓库的耦合模式变成了仓库统一管理、按需吞吐的中转模式。最关键的变化在于上游写完就完事了内存、磁盘立刻释放下游拉数据也不必再受制于上游节点是否还存活。3. 可靠性设计为什么敢把中间数据交给远程集群3.1 存储分层与两种核心模式要把 Shuffle 数据放到远程首先要回答放在哪、怎么保证不丢的问题。Uniffle 的 ShuffleServer 提供了两种存储模式对应配置项rss.storage.type的两个值MEMORY_LOCALFILE 模式刚接收的数据先写在 ShuffleServer 的内存里内存达到高水位阈值后再刷到 ShuffleServer 的本地磁盘。这个模式下可靠性来自多副本。理论上可以在服务器上配置多副本写入任何一个 ShuffleServer 挂了数据还能从另一个副本读取。适合不太依赖 HDFS、想自建轻量 Shuffle 集群的场景。MEMORY_HDFS 模式数据仍然先走内存缓冲但落盘目标是 HDFS。可靠性由 HDFS 的多副本机制保证适合已经有现成 HDFS 集群、希望把 Shuffle 数据也纳入统一存储管理的场景。两种模式下我建议优先看你们公司的底座。如果已经有 HDFS 且带宽充裕MEMORY_HDFS 运维上更省心——不用操心 ShuffleServer 本地磁盘坏了怎么办。如果没有 HDFS 或者不想让 Shuffle 数据占用 HDFS 空间MEMORY_LOCALFILE 加上副本配置也足够了。3.2 内存水位动态调整避免存储反噬计算ShuffleServer 本质上是一个存储进程最怕的就是内存被打满。Uniffle 对内存的管理有专门的高水位线和低水位线机制。当已使用的 Shuffle 内存超过高水位默认约 0.75~0.8存储层就开始强制刷盘当内存释放回落到低水位以下才停止刷盘动作。我看到很多初用者对为什么数据还要先在内存里待一会儿有疑问。答案是吞吐和合并。如果每来一批 64KB 的数据就直接落盘会产生大量小文件尤其 HDFS 模式下会生成海量小文件NameNode 压力巨大。先在内存攒一攒、按分区合并成大块再刷下去就是为了减少文件数量代价是需要 ShuffleServer 有足够内存作为缓冲带。这里有个实际调优心得ShuffleServer 的 JVM 堆外内存Direct Memory和堆内存都要认真规划不能光看堆大小。因为数据在内存缓冲阶段大量走的是 Direct Buffer堆外内存不够会比堆内存不够更容易出现诡异报错。3.3 故障场景下的降级与补偿再可靠的系统也要考虑万一。先看 ShuffleServer 宕机的场景。由于写入了多副本或在 HDFS 模式下有副本冗余下游读取时如果发现首选服务器连不上Coordinator 返回的服务器列表里还有备用副本客户端会切换节点继续拉取不会因为单点故障直接导致整个 Stage 失败。这比本地落盘模式下磁盘坏了就重算上游的体验好太多。再看 Coordinator 故障。Coordinator 是无状态的支持配置多个节点组成 Quorum客户端会轮询可用的 Coordinator。即使一个 Coordinator 挂了应用也不会中断只是某个 Shuffle 注册请求会稍微多一次重试。再看整个 ShuffleServer 集群容量不足的场景。Uniffle 的客户端在写入前会先从 Coordinator 获取服务器分配Coordinator 会基于内存、磁盘剩余空间等信息过滤掉不可用的服务器。如果你在部署时预留了缓冲容量这个机制可以大大减少写一半发现没地方存的尴尬。3.4 与 Executor 生命周期彻底解耦这是远程 Shuffle 最舒服的一点因为数据根本不在 Executor 本地所以 Executor 是可以随时安全地退出或重启的。在传统模式下动态分配Dynamic Allocation要想安全地缩容得依赖 External Shuffle Service 在节点上继续守护 Shuffle 数据。很多团队因为这个原因干脆不敢开动态分配或者把 Executor 的心跳调得很保守。用了 Uniffle 之后Executor 缩容只是释放计算资源Shuffle 数据安然躺在 ShuffleServer 上下游照样能拉到。这个解耦带来的灵活性和资源利用率提升比省下的那一点重算时间更有价值。4. 从零接入 Uniffle部署、配置与验证4.1 版本与兼容矩阵先瞄一眼再动手Uniffle 目前版本迭代很快接入前一定去官方文档确认兼容矩阵。以我写这篇文章时候的经验看主流的 0.8.x、0.9.x 版本对 Spark 2.4.x 到 3.5.x 都有官方支持Hadoop MapReduce 2.x/3.x 也能接Spark 是最成熟的一条路径。如果你的集群是 Spark 3.4 或更高版本还要留意 Uniffle 的 client jar 要和 Spark 版本匹配。一个常见错误是拿旧版 client 去跑新版 Spark结果 ShuffleManager 初始化时方法签名对不上直接报NoSuchMethodError。这个坑我在测试环境踩过一次排查了半天最后发现只是版本不匹配。4.2 部署一套最小 Shuffle 集群部署方式有 Docker、Kubernetes、Standalone 几种。这里给一个最精简的 Standalone 思路方便先跑通验证流程。假设你有三台机器分别是 shuffle01 / shuffle02 / shuffle03。部署内容如下在 shuffle01、shuffle02 上各启动一个 Coordinator 进程生产环境建议至少 2~3 个 Coordinator组成 Quorum在 shuffle01、shuffle02、shuffle03 上各启动一个 ShuffleServer 进程ShuffleServer 的核心配置示例rss.confrss.coordinator.quorum shuffle01:19999,shuffle02:19999 rss.storage.type MEMORY_HDFS rss.storage.basePath /rss/shuffle_data rss.server.buffer.capacity 2g rss.server.memory.shuffle.highWaterMark.ratio 0.75 rss.server.memory.shuffle.lowWaterMark.ratio 0.2 rss.server.port 19998 rss.hadoop.dfs.replication 2Coordinator 配置相对简单核心是把各 ShuffleServer 的地址和端口配好让注册流程能跑通rss.server.port 19998 rss.coordinator.rpc.port 19999启动顺序没有特别强制的要求但建议先起 Coordinator再起 ShuffleServer这样 ShuffleServer 一启动就能成功注册。启动后访问 Coordinator 的 Web 端口能看到当前注册的 ShuffleServer 列表、内存使用、磁盘剩余空间等信息这是后续排查问题最直接的入口。4.3 Spark 作业接入的完整配置把 Uniffle 的 client jar 放到 Spark 的 classpath 里可以用--jars引入然后在 Spark 配置里切换 ShuffleManagerspark-submit \ --master yarn \ --deploy-mode cluster \ --jars /path/to/uniffle-client-spark-xxx-shaded.jar \ --conf spark.shuffle.managerorg.apache.spark.shuffle.RssShuffleManager \ --conf spark.rss.coordinator.quorumshuffle01:19999,shuffle02:19999 \ --conf spark.rss.storage.typeMEMORY_HDFS \ --conf spark.serializerorg.apache.spark.serializer.KryoSerializer \ --conf spark.rss.client.read.buffer.size14m \ --conf spark.rss.writer.buffer.size3m \ --conf spark.shuffle.service.enabledfalse \ your-job.jar几个配置项简单解释一下spark.shuffle.service.enabledfalse建议显式关掉既然 Shuffle 数据已经不在节点本地再开 ESS 只会白白占用每个节点的资源spark.rss.client.read.buffer.size和spark.rss.writer.buffer.size分别控制读、写缓冲区的大小对 Shuffle 数据量极大的作业适当调大可以明显减少 RPC 次数但也要注意内存代价。另外强烈建议把序列化器设为 Kryo并注册好用户自定义类型。Uniffle 在传输层有独立的数据序列化处理Kryo 配合使用整体吞吐会比 Java 默认序列化好不少。如果你们的作业里有未注册的自定义类Kryo 会给出明确报错提前在测试环境过一遍就稳了。MapReduce 的接入思路类似核心是替换 Shuffle Consumer 插件、指定 Coordinator 地址具体配置方法以官方文档对应版本为准。这里不展开写具体类名因为不同版本略有差异照抄旧文章很容易踩坑。4.4 接入后怎么验证真的生效了配置完成后不要急着看性能。先做三件事验证链路是通的看 Spark 作业日志里有没有输出 Uniffle 相关的初始化信息比如加载了 RssShuffleManager、成功连接 Coordinator、为当前 Shuffle 分配到了哪些 ShuffleServer。打开 Coordinator 或 ShuffleServer 的监控页面观察作业运行期间是否有写入流量、当前正在服务的 Shuffle 数量、内存使用量。一个只读不写的作业监控上就应该动都不动。主动做一个杀 Executor实验在作业运行中手动 kill 掉一个 Executor观察 Spark 是否直接重算它的任务。如果用的是 Uniffle上游 Shuffle 数据已经在远程Kill 掉的 Executor 直拖着一个阶段的时间不会明显变长也不会有大范围重算。这三步走完基本可以确定你的作业真的在用远程 Shuffle 了。我见过有人配了一堆参数结果因为 client jar 没打入 Executor 的 classpathShuffleManager 悄悄回退成了普通模式性能一点没变还差点误判Uniffle 没用。所以验证生效这一步千万别省。5. 同场竞技Uniffle、ESS、Celeborn 到底怎么选5.1 三条技术路线的本质区别现在业界做 Shuffle 优化主流大致是三派保留节点本地模式的 External Shuffle ServiceESS、Uniffle 这种远程集中式 Shuffle、以及同样做远程 Shuffle 的 Apache Celeborn早期是阿里内部的 Remote Shuffle Service。把它们放在一起对比才看得清楚方案数据存放位置可靠性来源适用框架上手成本Spark 默认 ESSExecutor 所在节点本地磁盘ESS 进程独立于 Executor节点磁盘坏仍会丢数据仅 Spark几乎为零Apache Uniffle独立 ShuffleServer内存 HDFS/本地磁盘多副本 / HDFS 副本ShuffleServer 宕机可切换Spark、MapReduce 成熟其他框架在跟进需要部署 Coordinator ShuffleServerApache Celeborn独立 Shuffle 集群Master/Slave 模式双副本异步推送Spark、Flink 更友好需要部署 Master Worker从架构上看Uniffle 和 Celeborn 是很像的——都是把 Shuffle 计算和存储解耦但侧重点有差异。Uniffle 的优势在于成熟支持 MapReduce并且存储层能对接 HDFS这对于保留 HDFS 架构的企业来说很友好。Celeborn 在 Flink 场景的适配上走得比较靠前如果你的核心是 Flink 流计算 大量状态/窗口 Shuffle可能 Celeborn 更适合。5.2 我的选型判断这四种情况适合上 Uniffle结合我自己的生产经验遇到下面这些情况Uniffle 值得你认真评估Shuffle 数据量大的批处理作业尤其单作业 Shuffle 数据超过几百 GB磁盘和 IO 经常告警集群架构在逐步容器化Pod 本地盘不可控、不好扩容希望 Shuffle 数据落到独立存储集群作业稳定性波动频繁经常因为 Shuffle 重试导致下游延迟Kill Executor 成本高团队里同时有 Spark 和 MapReduce 作业想用一套组件统一 Shuffle 能力。反过来如果你们的作业普遍是 GB 级以下的轻量级任务或者核心场景是 Flink 流处理Uniffle 带来的收益就不明显引入一套新集群的运维成本反而可能大于收益。技术选型这事不是越先进越好而是匹配现状最划算的才是最好的。6. 生产环境落地时我踩过的几个实打实的坑6.1 ShuffleServer 的存储规划要当成状态服务来做ShuffleServer 是保存中间数据的它本质上是一个带状态的存储节点。很多人第一次部署时把它当普通无状态服务来规划随便挂一块云盘就上了。结果大数据量作业一来本地盘 IO 成为瓶颈Shuffle 数据写入变慢Map 端反过来被背压拖住。我的建议是MEMORY_LOCALFILE 模式下ShuffleServer 的磁盘一定要用高性能盘并且预留至少是日常峰值 Shuffle 数据量 1.5~2 倍的空间MEMORY_HDFS 模式下则要重点关注 ShuffleServer 与 HDFS 之间的网络带宽如果带宽被打满性能瓶颈就转移到了网卡上。6.2 不要开着默认配置就跑生产Uniffle 虽然开箱即用但默认配置是针对通用场景的。最容易出问题的两个默认项一个是highWaterMark.ratio内存富余的机器可以适当调高减少刷盘频次提升吞吐另一个是副本数像我前面说的MEMORY_LOCALFILE 模式一定要确认副本配置生效否则单个 ShuffleServer 宕机就是一场事故。还有个小细节ShuffleServer 的 JVM 堆大小和堆外内存要一起调。我用 16GB 堆 8GB 堆外跑中大批量作业GC 稳定吞吐也够。如果只调堆不管堆外数据缓冲段经常报 Direct buffer memory 溢出。6.3 先选一个作业灰度别一把梭我强烈建议不要第一天就把所有生产作业切到 Uniffle。合理的节奏是先在测试环境跑通验证链路再选一个 Shuffle 占比最高、业务价值最明确比如每天夜间跑批的作业做灰度和原有 ESS 模式对比运行时间、GC 时间、任务失败率、Yield 资源利用率等指标确认收益后再逐步推广。灰度期间的监控指标重点看两套一套是 Spark 作业自身的 Shuffle Read/Write 量、Shuffle 相关任务的失败重试次数另一套是 ShuffleServer 的内存水位、刷盘次数、RPC 响应延迟。这两套都正常才说明系统是健康的。6.4 小技巧利用 Coordinator 的监控做容量预测Coordinator 的监控页面能看到每个 ShuffleServer 的内存、磁盘、正在服务的 Shuffle 数量。我习惯每天早晨看一眼这些数据结合前一天的作业高峰基本能判断当前 Shuffle 集群的容量余量。这样在业务量上涨前可以提前扩容 ShuffleServer而不是等到作业快失败时才手忙脚乱。这套运维模式本质上已经不是在调一个组件而是在运营一个独立的存储服务了。想清楚这一点你对 Uniffle 的使用就会从接个插件升级为经营一个 Shuffle 存储平台上层稳定性也就真正有了保障。最后再分享一个扩展思路因为 Uniffle 把 Shuffle 数据集中化了后续你完全可以在上面叠加更精细的流量调度和缓存预热策略。比如让 Coordinator 根据历史热点提前把某些高竞争分区的数据调度到更近的 ShuffleServer 上。这些玩法在本地落盘时代想都不敢想但有了统一 Shuffle 引擎一切才刚起步。