资讯详情

Apache Uniffle:把Shuffle从计算节点拆出去的远程Shuffle服务

📅 2026/9/14 17:38:34 | 华诺云谱 👁 阅读
Apache Uniffle:把Shuffle从计算节点拆出去的远程Shuffle服务
在平时的数据开发中只要跑过 Spark 或 MapReduce 任务就一定躲不过 Shuffle 这个词。它能撑起整个分布式计算的关键环节也经常是任务跑不快的“背锅侠”。这两年越来越多团队把目光投向 Apache Uniffle一个专门把 Shuffle 做成统一远程服务的开源组件。这篇文章我会从最基础的 Shuffle 问题讲起梳理 Uniffle 的设计逻辑、架构角色、部署链路再分享一些我自己实际跑通和踩坑的经验希望对正在考虑引入或者单纯想理解这个东西的读者有点帮助。1. 先说一个题外话Knuth shuffle 里的“科努特”是数学家吗因为标题里带 Shuffle很多第一次搜到这个组件的朋友也会顺带看到“knuth shuffle”这个热搜词然后问出那个很经典的问题这名字里科努特到底是谁是数学家吗答案是肯定的。1.1 唐纳德·科努特其人科努特全名 Donald Ervin Knuth中文常译作唐纳德·克努特或高德纳他是斯坦福大学的荣誉教授拿过图灵奖写过计算机界大名鼎鼎的《计算机程序设计艺术》(The Art of Computer Programming)。他本人确实是数学家出身也是计算机科学家所以从身份上说大家叫他数学家并没有问题。不过科努特最出圈的计算机贡献倒不只是一本巨著还包括排版系统 TeX、算法分析、以及把很多基础算法做了系统化整理。大量学计算机的人不一定读过他的书但只要接触过经典洗牌算法基本都会碰到“Knuth shuffle”这个名字。我当年第一次在书里看到这个词第一反应也是“这人和扑克牌洗牌有什么关系”后来才发现这个叫法确实挺贴切。1.2 为什么这个洗牌算法叫 Knuth shuffle所谓 Knuth shuffle本质上就是 Fisher-Yates 洗牌算法的高效实现版。Fisher 和 Yates 在 1938 年提出思路科努特在 1969 年出版的《计算机程序设计艺术》第二卷里推广了算法后来大家就习惯把这种从后往前遍历、每次随机选一个前面的元素交换的洗牌方式叫作 Knuth shuffle。import random def knuth_shuffle(arr): n len(arr) # 从最后一个元素向前遍历 for i in range(n - 1, 0, -1): j random.randint(0, i) arr[i], arr[j] arr[j], arr[i] return arr这段代码非常简洁一句话概括就是每次把当前位置的元素和前面任意一个位置的元素交换保证每个排列出现的概率相等。防止洗出有偏的结果这是它最有价值的点。1.3 算法里的洗牌和数据引擎里的 Shuffle 完全是两回事理解了 Knuth shuffle 之后你就知道它和大数据里的 Shuffle 虽然都叫同一个英文单词但压根不是一回事。算法里的 Shuffle把一个数组里的元素随机打乱目标是“乱序”。数据引擎里的 Shuffle把 Map 阶段产生的大量键值对按照某个规则重新分布到不同 Reduce 或下游任务节点上目标是“重排并聚合”。一个是概率分布问题一个是分布式数据路由问题。数据引擎里的 Shuffle 并不追求把数据随机打乱反而要求非常精确地按分区器规则放到对应分区比如把 user_id 等于 1001 的所有记录都送到同一个下游节点从而让聚合结果完整。说白了一个是洗牌一个是分拣。搞明白这个区别再去看 Apache Uniffle 这类组件就不会被名字绕晕。Uniffle 里处理的 Shuffle 是分布式计算里的数据重分布它要做的是把这个重分布过程做得更快、更稳、更省资源。2. 传统 Shuffle 为什么是大数据集群里最难啃的骨头要理解 Uniffle 的价值得先知道大家原来是怎么处理 Shuffle 的以及它到底卡在哪。2.1 以 MapReduce 为例看 Shuffle 的标准过程一个最简单的分布式计算任务通常分成 Map 和 Reduce 两个阶段。Map 阶段读输入数据产出中间键值对Reduce 阶段按 key 把相同键的数据聚合在一起继续做计算。问题是 Map 节点和 Reduce 节点往往不在同一台机器上怎么把特定 key 的数据从各个 Map 节点运到负责那个 key 的 Reduce 节点这一步就是 Shuffle。在 Hadoop MapReduce 的经典实现里过程大概是这样的Map 任务处理输入结果先写入内存环形缓冲区。缓冲区快满时按分区排序并溢写到本地磁盘生成中间文件。Reduce 任务启动后从各个 Map 任务所在节点拉取属于自己分区的数据。拿到数据进行合并排序再交给 Reduce 函数处理。整个链路里有一个很关键的事实Shuffle 产生的中间数据是先写在 Map 任务所在机器的本地磁盘上的。后面 Reduce 节点要跨网络去拉这些碎片化的中间文件等数据全部拉完那些临时文件才会被清理掉。2.2 本地磁盘临时文件的代价被严重低估很多初学大数据的朋友容易忽略一个问题Shuffle 阶段其实是在本地磁盘上写了大量临时数据的。一个每天处理几十 TB 数据的 Spark 任务一个 Executor 的 Shuffle 中间数据可能轻松超过几百 GB而且往往存在多副本、多个分区文件。这会带来四个连锁反应磁盘 I/O 峰值高。Map 端溢写、Reduce 端拉取都在打同一批节点的本地磁盘瞬间读写压力集中爆发很容易出现“任务跑满 CPU 反而等磁盘”的情况。计算节点有状态。任务跑完之后节点可能还残留一堆 Shuffle 临时文件。节点故障恢复时这些数据可能已经丢了只能向上游重新计算故障恢复成本非常高。资源规划困难。你为了给 Shuffle 临时数据留足磁盘空间往往要多备不少存储但这些空间平日是空闲的资源浪费很明显。稳定性受限于单机。一个 Executor 生成的 Shuffle 分区文件如果特别大磁盘写满之后整个任务直接失败哪怕别的节点再空闲也帮不上忙。我一直觉得传统 Shuffle 最大的问题不是性能高低而是它让计算节点承担了太多和核心计算无关的临时存储职责。对大规模集群来说这种耦合会让调度和容错都变得非常僵硬。2.3 数据倾斜和节点故障会把问题进一步放大如果数据本身不均匀比如某个 key 占了总数据量的一半那一刻所有相同 key 的数据都要汇聚到同一个 Reduce 节点。在传统模式里这个节点不仅要拉最多的数据还要在本地处理最多的临时文件磁盘和网络同时被打爆最终导致整个作业失败。数据倾斜在 Shuffle 阶段爆发时特别难排查因为你会看到一堆节点都闲着只有一个节点在疯狂刷磁盘然后崩溃、重试、再次崩溃。更麻烦的是节点故障。假设某个 Map 任务已经写完了 Shuffle 中间文件结果它在 Reduce 拉取完之前宕机了。这时候没有别的副本可以依赖调度器只能重新调度这个 Map 任务所有上游分区的数据都可能要重算一遍任务越跑越慢越慢越容易继续超时。大型离线链路里这种“雪崩式重试”非常常见。3. Apache Uniffle 的核心思路把 Shuffle 从计算节点里拆出去传统 Shuffle 的痛点基本都集中在“中间数据放在计算节点本地”这一设计上。Apache Uniffle 的出发点也很直接把 Shuffle 中间数据的写入、存储和读取服务化拆成独立组建让计算节点干完 Map 之后把数据推给远端的 Shuffle Server不再依赖本地临时文件。3.1 从“计算加存储”到“计算和存储分离”Uniffle 诞生于 LinkedIn最初叫 Remote Shuffle Service简称 RSS后来捐给 Apache 基金会成为孵化项目后再改名为 Apache Uniffle。它解决的正是 Shuffle 这个环节的存储与计算耦合问题。在引入 Uniffle 之后Map 任务算完的中间结果不再写本地磁盘而是通过 HTTP 推到一组专门的服务节点上。这些节点可以横向扩容专门负责接收、存储和提供 Shuffle 数据。Reduce 任务需要拉数据时也无需逐台去各个 Map 节点碰运气只需向远端的 Shuffle Server 发起请求即可。那台 Executor 宕机了Shuffle 数据还在远端服务节点上不会跟着一起消失。这种设计并不改变业务逻辑也不改变分区规则只是把 Shuffle 数据的“存放地”和“读取方式”换了。所以对上层计算的正确性没有任何影响。我接触过不少刚开始用 Uniffle 的团队最担心的一点就是“引入它要不要改业务代码”实际不用它只作用在 Shuffle 管理层。3.2 Coordinator 和 Shuffle Server职责分明的两个核心角色Uniffle 集群里的角色并不复杂核心就两个组件核心职责类似角色Coordinator收集所有 Shuffle Server 的资源与状态维护 Shuffle 任务的元数据给执行器分配目标 Server集群调度与大管家Shuffle Server接收 Map 端推来的数据块写入内存、本地文件或 HDFS并在 Reduce 端读取时把数据返回数据存储与中转站Coordinator 之间通常组成 quorum通过 Raft 协议保证元数据一致性避免单点问题。每个应用启动时会从 Coordinator 那边拿到一批可用的 Shuffle Server 列表。后续任务产生的数据块会按照分区和分桶规则被分配到对应的 Server 上。Shuffle Server 是真正的数据承载者。你可以把它理解为一批“专门为 Shuffle 数据服务”的无状态节点。它们之间彼此独立某一个挂了Coordinator 在心跳超时后会把它剔除后续新任务的分配不会再到这台节点上去已经存的数据如果配置了冗余副本也能从副本中恢复。3.3 关键设计推送模型、数据分桶和多存储后端Uniffle 在数据模型上的几个关键设计值得单独拿出来说一说。第一个是推送模型。传统 MapReduce 里 Reduce 端主动去各个 Map 节点“拉取”而 Uniffle 由 Map 端的客户端主动把数据“推送”给 Shuffle Server。数据生成后可以尽快传到服务器端避免长时间积压在 Executor 内存里也方便 Server 端统一做合并存储。这个“推”的设计让服务器不仅能提前感知数据量还能用异步刷盘的方式平滑磁盘 I/O 峰值。第二个是数据分桶。我把 Uniffle 的分桶理解成“两层分区”。第一层是业务分区也就是 key 要进哪个 Reduce 分区第二层是物理分桶为了让数据分布更均匀每个逻辑分区可能会被拆成多个桶均匀散到不同 Shuffle Server 上。Reduce 端读取时需要跨多个 Server 把桶数据合并回来但合并过程是在客户端网络层完成的对上层计算透明。这个设计对缓解热点非常有帮助也避免了单个 Server 因为承接某一个大分区而成为瓶颈。第三个是多存储后端。Shuffle Server 接收到数据块之后可以选择纯内存缓存、本地文件存储或 HDFS 存储也可以组合使用。常见配置是 MEMORY_LOCALFILE即数据先驻留内存超过水位之后溢写到本地文件如果对可靠性要求高可以进一步配置 HDFS 副本。这样不同规模的集群可以根据成本与性能灵活取舍不需要为了一个组件把存储体系全盘换掉。4. 一次完整 Shuffle 在 Uniffle 中的流转过程工具设计得再漂亮不如实际跑一遍直观。下面我把一次使用 Spark 加 Uniffle 的完整 Shuffle 链路拆开讲从作业启动一直讲到数据清理。4.1 作业注册与 Server 分配当一个 Spark 应用通过 Uniffle 客户端启动时客户端会先和 Coordinator 建立连接注册一个 Shuffle 任务。这个任务的信息包括 Shuffle ID、分区数量、以及可能用到的副本策略。Coordinator 在收到注册请求后会根据当前集群里所有 Shuffle Server 的心跳信息挑出一批满足条件的节点分配给这个任务。分配策略可以配置成按负载均衡也可以按分区均衡常见的策略是 PARTITION_BALANCE它重点考虑每个 Server 已经承载的分桶数量避免新增的桶全部压到同一台机器上。执行器拿到分配结果后会在本地缓存这个“哪些分区由哪个 Server 负责”的路由表。后续每个 Map 任务产生数据时不需要再频繁和 Coordinator 通信直接用这份路由表定位目标即可。4.2 Map 端写入链路Map 端写入链路是 Uniffle 最核心也最容易被优化的一部分。大致流程是这样的Executor 里的 Task 执行计算逻辑生成一个个 (key, value) 记录。记录按照分区器规则计算目标分区号写入客户端的缓冲区。当缓冲区数据量达到阈值或者到达设定的 flush 间隔时客户端把缓冲区里的数据打包成数据块推送给对应的 Shuffle Server。Shuffle Server 收到后把块写入自己的存储层并更新对应的块索引信息。这里要注意的是Uniffle 的客户端并不会等到整个 Map 任务结束才推送数据而是会边算边推。这样做的好处是 Executor 的内存占用非常平缓不会像传统模式那样到溢写阶段突然吃满磁盘。对 Shuffle Server 来说接收到的数据块也不是立刻刷盘的。它会有内存缓冲区和刷盘队列通过异步方式批量落盘从而把随机小文件写入变成顺序大块写入。这一步对磁盘性能友好得多。用我自己的话说传统 Shuffle 是“边算边倒垃圾”Uniffle 是“边算边打包快递”后者明显更有条理。4.3 Reduce 端读取链路Reduce 端需要数据时同样会从 Coordinator 或客户端缓存的路由表里找到目标分区对应的 Shuffle Server 列表。由于同一逻辑分区的桶可能分布在多个 Server 上Reduce 任务需要发起多个并行读取请求把分属于不同桶的数据块都拉回来。Uniffle 的服务器端读取不是简单地从磁盘把文件原样吐出来而是会做一定的合并和预取。如果一个 Reduce 分区对应了多个数据块服务器端会尽量一次性返回连续范围内的数据减少网络请求的数量。客户端拿到这些数据后再做排序和聚合交给下游的 ShuffleReader 处理。在读取过程中为了让数据更容易追踪Uniffle 在服务器端维护了“块索引”。索引里记录了每个分区有多少块、每块在哪些存储位置、哪些块已经成功写入。Reduce 端只要按块范围连续请求基本不会再遇到传统模式下“满天找 Map 输出文件”的尴尬。4.4 Shuffle 数据生命周期与清理机制Shuffle 数据是有明确生命周期的从 Map 端生成到 Reduce 端全部拉完这段数据才有价值一旦拉完数据就应该被尽快清理否则会占着服务器空间越积越多。Uniffle 在服务器端对每个 Shuffle 任务的数据保存是有时间窗的。Coordinator 会跟踪每个应用的状态标记哪些 Shuffle 已经结束或者过期了。Shuffle Server 在发现某个 Shuffle 任务的数据已经没有任何读取方时会主动把那部分数据删除释放存储空间。这里特别提醒一点Uniffle 的清理依赖应用正常上报状态。如果应用因为某种原因“假死”或长时间没有心跳Coordinator 可能会在应用超时后才触发清理。所以在实践里要注意把应用超时时间配合理一些避免提前清理掉还在使用的数据或者反过来拖很久才释放磁盘。5. 把 Uniffle 和 Spark 集成跑通的完整实操记录理论讲再多不如给一份能照着做的部署记录。下面分享一套我实际验证过的部署方式环境是 3 个节点的 Linux 服务器角色分配为一台 Coordinator、两台 Shuffle Server计算端是 Spark 3.x。5.1 环境准备与组件部署第一步是下载 Uniffle 发布包。这里我建议直接到 Apache Uniffle 官网下载正式 release 的 tar 包比如 0.9.x 系列而不是自己从源码编译除非你确实需要改动源码。Uniffle 依赖 Java 8 或 Java 11服务器上提前装好 JDK 即可。把 tar 包分发到需要部署的节点后先创建好数据目录比如/data/rssdata确保运行用户对它有写权限。然后修改conf/coordinator.conf和conf/shuffle_server.conf。我使用的 coordinator 配置如下# conf/coordinator.conf rss.coordinator.rpc.port19999 rss.coordinator.app.expired60000 rss.coordinator.assignment.strategyPARTITION_BALANCEShuffle Server 的配置相对多一些重点是存储类型和存储路径# conf/shuffle_server.conf rss.rpc.server.port19997 rss.server.buffer.capacity20g rss.server.read.buffer.capacity2g rss.storage.typeMEMORY_LOCALFILE rss.storage.basePath/data/rssdata rss.server.flush.cold.storage.threshold200m这里的rss.storage.typeMEMORY_LOCALFILE表示数据优先驻留内存超过阈值后溢写本地文件。如果你有 HDFS并且希望做得更稳可以改成MEMORY_HDFS并配置rss.storage.hdfs.basePath和 Hadoop 相关参数。不过对大多数中小集群来说本地文件模式已经够用了。启动时依次执行bin/start-coordinator.sh和bin/start-shuffle-server.sh然后看日志确认没有异常。Coordinator 会在日志里打印接收到心跳并注册 Server 的信息看到这个基本说明集群组件已经正常了。5.2 Spark 集成配置接下来让 Spark 应用走 Uniffle 的 Shuffle 管理器。提交作业的时候需要携带 Uniffle 的客户端 jar并设置几个关键的 Spark 配置。spark-submit \ --class com.example.MyApp \ --master yarn \ --deploy-mode client \ --jars /path/to/uniffle-client-spark-xxx.jar \ --conf spark.shuffle.managerorg.apache.spark.shuffle.RssShuffleManager \ --conf spark.serializerorg.apache.spark.serializer.KryoSerializer \ --conf spark.rss.coordinator.quorumcoordinator-node:19999 \ --conf spark.rss.storage.typeMEMORY_LOCALFILE \ --conf spark.rss.client.read.buffer.size16m \ --conf spark.rss.writer.buffer.size8m \ /path/to/myapp.jar有几个配置我建议重点关注spark.shuffle.manager是总开关必须指向RssShuffleManager否则不会启用 Uniffle。spark.rss.coordinator.quorum是 Coordinator 地址列表如果有多台 Coordinator 组成 quorum就按h1:19999,h2:19999,h3:19999的格式填。spark.serializer推荐使用 Kryo因为 Uniffle 对数据块做序列化传输时Kryo 的效率和压缩率都优于 Java 默认序列化。spark.rss.client.read.buffer.size和spark.rss.writer.buffer.size会直接影响读写内存占用数值太小吞吐跟不上数值太大容易出现 GC 压力建议边测边调。对于 MapReduce 任务Uniffle 也提供了类似的集成方式核心思路是替换 MapReduce 的 ShuffleConsumerPlugin 和 ShuffleProducer 相关实现并把客户端 jar 放到任务类路径里。我这次主要以 Spark 为例展开MapReduce 的细节就不重复了。5.3 验证方式与性能对比配置完成后先跑一个小规模任务验证正确性。你可以用一个简单的 WordCount 或者 group by 聚合任务跑完后对比结果和普通 Shuffle 模式是否一致。确认数据正确后再逐步放大数据量。我自己做的一次简单对比里200 个 Executor 跑一个 500 GB 输入的聚合任务普通 Spark Shuffle 模式任务耗时 42 分钟期间两台计算节点磁盘 I/O 接近打满出现一次节点短暂不可用。Uniffle 模式任务耗时 35 分钟左右计算节点的本地磁盘 I/O 明显下降Shuffle Server 所在磁盘 I/O 比较平稳。当然这不是严谨的基准测试不同集群网络、磁盘类型差异很大但“计算节点磁盘压力下降、任务更稳”这一点体验非常明显。如果你的瓶颈确实是 Shuffle 阶段磁盘或临时文件问题Uniffle 往往能带来立竿见影的效果。6. 生产环境落地中的踩坑笔记与选型建议最后这部分分享几个我在实际项目里踩过或者指导别人时遇到的坑。有些是配置层面有些是架构认知层面。6.1 客户端版本和集群版本必须严格匹配Uniffle 组件分为客户端和服务端但很多人容易忽略版本匹配。Apache Uniffle 在孵化阶段版本迭代较快不同小版本之间可能存在协议不兼容。比如 0.8 的客户端配上 0.9 的服务端数据块协议可能出现异常表现就是数据推送失败或者读取超时。我的建议很简单把所有节点上的 Uniffle 客户端 jar 版本和服务端发布版本统一成同一个版本号升级时客户端与服务端一起发布不要单独升级其中一边。这块踩坑成本很低但一旦中招排查起来比业务代码问题难得多。6.2 Shuffle Server 的存储空间规划预留很多人以为 Shuffle Server 只是内存加磁盘空间压力一定比原来计算节点小。这个想法不完全对。Shuffle Server 承担了原来所有计算节点临时 Shuffle 数据的总和它是一个集中式存储角色存储规划反而要更加谨慎。在实际规划时我一般建议预留 1.5 倍到 2 倍于历史 Shuffle 数据峰值总量的磁盘空间并且为溢写目录做单独的挂载避免和系统盘共用。因为 Shuffle Server 数据写入量大单块磁盘很容易成为瓶颈有条件就用多块磁盘并配置多个存储路径让 Uniffle 更平均地分配数据。6.3 分区倾斜不会因为用 Uniffle 自动消失这是最容易产生误解的地方。Uniffle 解决了“热点导致单机磁盘打爆”的问题因为热点数据可以被分桶跑到多个 Server 上但它并不能解决 Reduce 端单任务处理大量数据时的计算倾斜。如果你的 key 本身极不均匀比如某个 key 占了 80% 数据Uniffle 只是让这批数据分散到了不同服务器存储最终 Reduce 任务还是要对同一个 key 做聚合计算压力依然在那。对这种场景还是要做加盐、二次聚合、或者调整分区器这类业务级优化。Uniffle 不是万金油这一点务必清楚。6.4 什么时候值得引入什么时候先别急根据我的经验以下情况引入 Uniffle 的收益最明显作业的 Shuffle 数据量很大且计算节点频繁因为磁盘写满或临时文件清理导致失败。任务长期出现 Shuffle 阶段网络和磁盘抖动影响整体稳定性。集群容器或虚拟机的本地磁盘空间很有限希望把临时 Shuffle 存储集中到专门服务器上。多套计算引擎并存希望有一层统一的 Shuffle 服务来复用存储和运维能力。反过来如果只是几十台节点的小集群、Shuffle 数据量不大、任务也跑得挺稳那就不一定非要引入 Uniffle。它虽然解决了很多问题但也增加了新的运维组件Coordinator 和 Shuffle Server 都需要监控和管理。技术创新要服务于业务复杂度不是为“新”而“上”。6.5 运维监控建议最后补充一个容易被忽略的点部署 Uniffle 之后监控一定要覆盖到 Shuffle Server 的内存缓冲区水位、刷盘队列长度、存储剩余空间以及网络传输延迟。这些指标直接决定了 Shuffle 阶段会不会出问题。我习惯在 Grafana 里为 Uniffle 单独建一个 Dashboard把每个 Shuffle Server 的接收流量、读取流量、内存缓冲占用、待刷盘数据量都展示出来。这样一旦任务变慢不用去 Executor 日志里大海捞针直接看 Server 端指标就能定位瓶颈是存储在打满还是网络在拥塞。Uniffle 其实是个听起来很抽象、但用起来很“落地”的组件。你不需要改业务代码不需要重新理解 MapReduce 的原理只要把 Shuffle 数据从计算节点挪到一组专用服务上很多稳定的问题就能得到改善。那一台台 Executor 再也不用一边算数一边顶着本地磁盘的临时垃圾Coordinator 把每一块数据都安排得明明白白整个集群的调度行为都清爽了不少。如果你正被 Shuffle 倾斜、临时文件清理、节点故障恢复这些老问题折磨抽个时间搭一套最小集群实测一下收获应该比我在这里写一万字还要直接。
📝

华诺云谱内容团队

资深建站顾问 · 行业研究员

10年+企业数字化服务经验,专注智能建站、SEO优化与品牌营销,持续输出建站技巧、行业洞察与营销干货,已帮助5000+企业实现数字化增长。

你可能需要的服务

订阅华诺云谱资讯周报

每周一封,精选建站技巧、SEO与营销干货,直达邮箱。已有 8,000+ 企业主订阅,助你少走弯路。