搞大数据的,几乎没有人不被 Shuffle 折磨过。尤其是当单批次作业处理的数据量冲到 PB 级时,Shuffle 阶段不仅仅是慢,而是直接能把集群磁盘打满、把网络打瘫、把任务拖到超时。vivo 的大数据平台在走向 PB 级规模的过程中,同样踩遍了这些坑,最终选择了 Apache Celeborn(Remote Shuffle Service,RSS)作为核心优化方案,逐步把 Shuffle 这块硬骨头啃了下来。这篇文章就围绕我们在这条路上做过的选型对比、架构改造、参数调优和问题排查展开,内容全部来自真实落地场景,希望对正在被 Shuffle 困扰的团队有直接的参考价值。
很多人第一次听到“Shuffle”,脑子里蹦出来的是算法课里那个给数组随机打乱的 Knuth Shuffle,顺便还会好奇科努特到底是不是个数学家。实际上分布式计算里的 Shuffle 完全是另一码事——它不是洗牌,而是把数据按照 key 重新分组、跨节点传输的过程。Map 端产出的中间结果要发给对应的 Reduce 端,这个“发”的过程就是 Shuffle。数据量小的时候不觉得,一旦到了 PB 级,Shuffle 的代价会被放大到让人绝望。
1. 先搞明白:Shuffle 为什么会成为性能杀手
1.1 从一次全网用户行为分析任务说起
我们平台上有大量用户行为日志分析作业,每天处理的数据量在几百 TB 到几 PB 之间。这类作业的典型特征是:Map 阶段吞吐很高,每个节点每秒能处理几万到几十万条记录,但一到 Shuffle 阶段就原形毕露——Reduce 端要等 Map 端把所有中间结果写完、再跨节点拉取,链路长、环节多,任何一个环节卡住都会拖慢整个作业。
当时我们遇到一个非常典型的问题:一次涉及 5000 多个 Map 任务、2000 多个 Reduce 任务的分析作业,Map 阶段只跑了 12 分钟,Shuffle 阶段却用了 47 分钟。而且 Shuffle 期间,部分节点的磁盘 IO 长时间处于 100% 饱和状态,网络峰值带宽也被占满,直接导致同一批跑在集群上的其他作业也跟着遭殃。
1.2 Shuffle 慢的三个根源
总结下来,大规模 Shuffle 的性能瓶颈主要来自三个方面。
第一,小文件问题。这是最经典的痛点。Spark 默认的 Hash Shuffle 在 Map 端会为每个 Reduce 分区生成一个文件,m 个 Map 任务、r 个 Reduce 分区就要产生 m×r 个文件。5000×2000 就是一千万个文件,光是文件元数据的管理和清理就够 NameNode 喝一壶的。虽然后续版本引入了 Sort Shuffle 和 Consolidation 机制来减少文件数量,但本质上 Map 端还是要写大量中间文件,磁盘 IO 压力依然很大。
第二,网络传输放大。Shuffle 本质上是一个全对全(All-to-All)的数据交换过程,数据量越大,网络传输的时间占比就越高。尤其是当数据分布不均匀时,某些节点需要拉取的数据量远超平均值,形成网络热点。
第三,Failover 成本高。这一点最容易被低估。Map 端写好的中间数据存放在本地磁盘,一旦某个节点宕机,所有依赖该节点数据的 Reduce 任务都得重新计算,整个作业可能从头再来。在 PB 级数据场景下,这种重算成本是灾难性的。
所以说,Shuffle 优化不能只盯着某个参数调一调,必须从架构层面解决问题。这也是我们后来坚定选择 Celeborn 的根本原因。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 为何选 Celeborn:基于 RSS 思想的架构革新
2.1 传统 Shuffle 方案的局限
在决定引入 Celeborn 之前,我们先后评估过几条路径。
一是继续用 Spark 自带的 Sort Shuffle。它确实比 Hash Shuffle 进步了很多,但在超大规模场景下,Map 端写本地磁盘、Reduce 端跨节点拉取的模式没有变,磁盘故障导致的重算风险也没有消除。二是引入 External Shuffle Service(ESS),这个方案解决了 Executor 与 Block 生命周期分离的问题,但中间文件依然在本地,小文件问题没有本质改善。三是尝试过自研 Shuffle 服务,但考虑到底层存储、高可用、数据一致性这些环节的复杂度,团队评估后觉得投入产出比不高。
2.2 Celeborn 的核心理念:把 Shuffle 搬到独立集群
Celeborn 走的是另一条路——把 Shuffle 数据从计算节点本地搬到一组独立的、专门部署的 Shuffle 服务集群上。这个思路听起来简单,但带来的收益是全方位的。
首先,Map 端不再往本地写数据,而是通过异步方式把中间结果推送到 Celeborn 集群。Reduce 端直接从 Celeborn 拉取数据,中间的传输链路从“Map 本地 → Reduce”变成了“Map → Celeborn → Reduce”,看似多了一跳,但因为 Celeborn 集群的存储和网络都是专门为 Shuffle 场景设计和调优的,整体性能反而更高。
其次,Celeborn 在接收端做了数据合并。它会把多个 Map 任务产生的、属于同一个 Partition 的数据块在服务端合并成一个大文件,从根本上解决小文件问题。我们在生产环境实测过,引入 Celeborn 之后,NameNode 上 Shuffle 中间文件的数量下降了 90% 以上,NameNode 的 RPC 压力明显缓解。
第三,Celeborn 支持数据的副本机制。默认情况下每个数据块会写入两个 Worker 节点,任何一个节点宕机,另一个节点上的副本可以直接接管,Reduce 端无需等待重新计算。这个特性在 PB 级作业调度时太重要了——大作业跑十几个小时,中间任何一个节点故障都可能导致前功尽弃,有了副本机制之后,故障恢复时间从小时级降到了分钟级。
2.3 vivo 为什么敢在 PB 级场景用 Celeborn
选型的时候,我们也不是没有顾虑。当时 Celeborn(那时还叫 Remote Shuffle Service)在社区的知名度远不如现在,生产环境的成功案例也不多。但我们在深入看代码、做了小规模压测之后,判断这个方向的架构设计是对的,而且社区迭代速度很快。最关键的一点是,Celeborn 对上层计算引擎透明——升级 Celeborn 不需要修改 Spark 或 Flink 的作业代码,只需要调整配置并重启作业。这个特性大大降低了试错成本,即使效果不达预期,回滚的代价也很小。
我们决定先在几条业务线的小作业上跑通,再逐步扩大到中大型作业,最后才在 PB 级作业上全面开启。整个过程遵循“灰度验证、逐步放量”的原则,没有一上来就搞“大跃进”。
3. 落地实操:从部署到核心参数调优
3.1 部署架构与前置依赖
我们集群的 Celeborn 部署架构并不复杂,核心组件就三个:Master、Worker、客户端(内嵌在 Spark/Flink 应用中)。
Master 负责管理 Worker 节点、维护 Partition 与 Worker 的映射关系,并对外提供 Shuffle 状态的查询接口。为了保证高可用,我们部署了三个 Master 节点,通过 ZooKeeper 做选主和状态同步。Worker 是真正干活的节点,负责接收 Map 端推送的数据、写入磁盘、响应 Reduce 端的拉取请求。每个 Worker 节点分配了 64GB 内存、8 块 NVMe SSD,数据落地前先缓存在内存中,达到阈值后再异步刷盘。
这里有个容易被忽略的点:Celeborn 的存储目录不要和计算集群的本地目录混用。我们单独挂载了专用的数据盘,并且把 celeborn.worker.storage.dirs 配置为多个磁盘目录以均衡 IO。如果复用 Hadoop 的数据目录,容易出现磁盘空间被 Shuffle 数据占满、影响其他组件的情况。
我们在生产环境使用的版本是 Celeborn 0.3.x,Spark 3.3、Flink 1.16 都做过了完整适配。如果你所在的公司还在用 Spark 2.x,建议先确认 Celeborn 是否支持对应版本,避免适配成本超出预期。
3.2 关键配置参数逐项解析
Celeborn 的参数并不多,但每个参数都值得细调。我用表格梳理一下我们线上实际生效的核心配置:
| 参数项 | 推荐值 | 说明 |
|---|---|---|
| celeborn.worker.storage.dirs | 8 个 NVMe 目录 | 按需配置,注意各盘 IO 能力对齐 |
| celeborn.worker.memory | 48GB ~ 64GB | 视单机可用内存而定,内存太小会导致刷盘过于频繁 |
| celeborn.push.maxReplicate | 2 | 数据副本数,PB 级作业建议至少 2 副本 |
| celeborn.client.push.buffer.size | 64KB ~ 256KB | Map 端推送缓冲区大小,影响推送吞吐 |
| celeborn.worker.flush.buffer.size | 256KB | 刷盘缓冲区,过小会导致磁盘写入次数增加 |
| celeborn.worker.disk.flusher.threads | 4 | 刷盘线程数,NVMe 盘可以适当调大 |
| celeborn.client.fetch.maxRetries | 3 | 拉取失败重试次数,防止偶发网络抖动 |
| celeborn.worker.partitionSorter.enabled | true | 开启分区排序,提升拉取效率 |
这里面最值得多说两句的是 celeborn.push.maxReplicate。如果你对数据安全性要求极高,可以设成 3,但要注意磁盘占用会翻三倍。我们对不同作业做过对比,2 副本和 3 副本在性能上差异不大,但磁盘成本差了 50%,所以最终定了 2 副本。如果你的作业主要是短小低频的分析任务,1 副本来跑也未尝不可,但前提是你能容忍极端情况下任务重算。
3.3 PB 级数据量的任务调优组合
以我们线上一个日均处理数据量约 1.2PB 的用户标签计算作业为例,任务规模大概是:Map 端 6000 个并发、Reduce 端 3000 个并发,每个任务处理数据量在 200MB 左右。
资源配置方面,我们给每个 Executor 分配了 4 核 8GB 内存,Map 端的内存里单独留了 spark.shuffle.push.buffer.size 的空间给 Celeborn 客户端做推送缓冲,这个参数我们设成了 128KB。一开始用了默认的 32KB,结果推送次数过于频繁,CPU 上去了但吞吐没上来,调大之后效果立竿见影。
我们还在 Spark 侧做了两个关键切换:一是把 spark.shuffle.manager 换成 org.apache.spark.shuffle.celeborn.RssShuffleManager,二是把 spark.serializer 保持为 Kryo 不变。Celeborn 对 Kryo 支持得很好,不需要额外修改序列化方式。
对于 Flink 作业,我们的做法是在 flink-conf.yaml 里添加 shuffle-service-factory.class: org.apache.celeborn.client.flink.CelebornShuffleServiceFactory,同时把 execution.shuffle-mode 调成 ALL_EXCHANGES 之外的模式以适配 Celeborn 的推拉机制。Flink 的适配比 Spark 稍微复杂一些,建议先在测试环境跑一个简单的流式 WordCount 验证链路通不通,再接入真实作业。
3.4 集成细节与踩坑提醒
集成 Celeborn 时最容易踩的坑是版本不一致。客户端的 jar 包版本必须和集群的 Master/Worker 版本保持一致,否则会出现协议不兼容、数据推不上去的诡异问题。我们在测试环境就因为客户端用了 0.3.0 的 jar、集群跑的是 0.3.2,结果出现间歇性的 Push 超时,排查了很久才定位到是版本问题。后来我们干脆把所有依赖 Celeborn 的作业镜像统一打了固定的 jar 包版本,才从根本上杜绝了这个问题。
另一个容易被坑的点是动态资源分配。如果你在 Spark 作业里开了 Dynamic Allocation,Executor 的个数会动态变化。Celeborn 官方文档明确要求在启用 RSS 时关闭动态资源分配,因为 Executor 动态增减会导致 Shuffle 数据分布失衡,甚至出现数据丢失。我们一开始没注意,线上就出现过 Executor 缩容后部分 Partition 数据找不到的问题。关闭动态资源分配之后,一切回归正常。
4. 性能优化三板斧:合并、压缩、降 IO
4.1 小文件合并机制是如何工作的
Celeborn 在服务端对数据做了分区级的合并。每个 Partition 不再对应一堆小文件,而是维护一个不断追加写入的大文件——严格来说是文件组。Worker 端会按照 Partition ID 建立独立的目录和数据文件,Map 端推过来的数据块先缓存在内存里,由 Flusher 线程异步写入磁盘,写入时会对同一 Partition 的数据做聚合追加。
这个设计的好处在于,Reduce 端拉取数据时,只需要从 Worker 上顺序读取一个大文件,减少了大量随机 IO。在我们实测中,开启 Celeborn 之后,单批作业的 Shuffle 读耗时从原来的 18 分钟降到了 9 分钟左右,是肉眼可见的提升。对于整个集群来说,NameNode 的文件数量压力也大幅下降。
4.2 压缩算法与编解码器选型
Shuffle 数据在网络中传输时,压缩是降低带宽占用的最直接手段。Celeborn 默认使用 LZ4 压缩,但我们一开始发现 LZ4 在低压缩比的场景下收益不太明显。后来对作业的数据特征做了分析,发现大批量字符串类型的数据,用 ZSTD 压缩能获得更高的压缩比,CPU 开销增加大约 8%,但网络传输量下降了近 35%。
所以我们的做法是:默认压缩算法保持 LZ4 不变,对于网络带宽紧张、作业对延迟不敏感的批处理任务,在 Spark 客户端通过 spark.shuffle.push.compressAlgorithm 参数切换为 ZSTD。这里有一个细节需要注意,如果启用了 Celeborn,还要单独配置服务端的 celeborn.shuffle.compression.codec,两边必须保持一致。如果压缩算法不一致,Reduce 端拉到的数据就无法解码,作业会报非常诡异的异常。
4.3 推拉模式优化与流量控制
Celeborn 同时支持 Map 端主动推送和 Reduce 端主动拉取,这两种模式的衔接直接影响整体性能。
早期版本里,所有 Partition 的数据都通过同一个网络连接推送到 Worker,在并发量大时容易造成 TCP 窗口拥塞。我们在新版本中启用了 celeborn.client.push.connectionPerPartition,允许每个 Partition 单独建连。这个参数打开后,并发推送的效率提升非常明显,尤其是在多个 Map 任务同时处理同一个大表时,不会再出现互相争抢同一个连接的情况。
Reduce 端的拉取也要注意控制并发度。spark.reducer.maxSizeInFlight 这个参数设得太小会导致拉取速度上不去,设得太大又会造成 Worker 端内存压力陡增。我们试过 48MB、64MB、96MB 三档,最终定在 64MB,整体效果最好,既能把网络带宽用满,又不会导致频繁的 GC。
5. 稳定性保障:PB 级场景的可用性设计
5.1 Master 高可用与故障切换
Celeborn 的 Master 通过 ZooKeeper 做选主,当 Active Master 宕机后,Standby Master 会自动切换为 Active。切换过程对上层作业基本透明,但需要注意一点:Master 切换后,部分正在进行中的 Push 操作会有短暂的重试窗口。我们在测试时多次 kill 掉 Active Master,观察到的现象是作业会卡住十几秒,之后自动恢复,没有出现过数据丢失。
如果你不想依赖 ZooKeeper,Celeborn 也支持基于 Raft 的高可用模式。我们没有使用 Raft,主要是运维习惯的问题,团队对 ZooKeeper 的运维经验更丰富。两种模式各有优劣,建议根据团队实际情况决定,而不是盲目跟风。
5.2 数据一致性:ACK 与重试机制
Push 数据的一致性通过 ACK 机制保证。Map 端每推送一个数据块,Worker 成功写入并刷盘后返回 ACK;Map 端收到 ACK 才认为该数据块推送成功。如果 Worker 返回异常或者长时间没有响应,Map 端会重试推送。
这个机制有一个潜在问题:如果 Worker 在刷盘完成前宕机,但数据块的副本已经写到了第二个 Worker 上,此时第一个 Worker 的 ACK 没有返回,Map 端会重新推送。重推时如果第二个 Worker 已经持有相同的数据块,会不会造成重复存储?我们带着这个疑问查了源码,发现 Celeborn 在数据块中带了一层序列号(batchId)去重。除非你自己改造了客户端逻辑,否则默认实现不会造成重复数据问题。
在实际生产中,网络抖动导致 Push 重试是常态。我们把客户端侧的重试间隔从默认的 5 秒调低到了 2 秒,重试次数从 3 次调到了 5 次。不要小看这两个参数,在集群网络有偶发拥塞的时段,适当调低间隔、增加次数,能避免大量任务因为一次网络抖动集体失败。
5.3 限流与熔断保护
Shuffle 数据集群如果没有任何保护机制,一旦有多个超大作业同时运行,很容易把 Worker 的内存和磁盘打满。我们做了两层保护。
第一层是内存限制。Worker 进程的内存使用量超过 celeborn.worker.memory.monitorThreshold 后,会触发 Spill 操作,把内存中的数据提前刷到磁盘。这个阈值我们设成了 0.8,也就是内存使用达到 80% 就触发刷盘。因为 Worker 的缓存数据本身是异步刷盘的,提前 Spill 对性能的影响很小,但能避免 OOM。
第二层是磁盘保护。Celeborn 会定期检查磁盘可用空间,如果某个目录的剩余空间低于阈值,会停止向该目录写入新的数据块,并把数据重定向到其他目录。我们遇到过一次 Worker 磁盘被日志文件占满的情况,就是靠这个机制保住了 Shuffle 数据盘,不然大量作业都得挂掉。
6. 常见问题与排查经验实录
6.1 数据倾斜导致单 Worker 热点
引入 Celeborn 之后,我们遇到过一类新问题:某个 Worker 节点的 IO 特别高,但其他 Worker 节点很空闲。排查下来发现,是部分作业的数据倾斜导致的——某个 Partition 的数据量远超其他 Partition,而 Celeborn 按照 Partition 粒度分配存储,热点 Partition 所在的 Worker 自然就成了瓶颈。
这个问题没有彻底根治的手段,只能缓解。我们的方案是让业务侧在写入时按照更细粒度的 key 做预聚合,减少数据倾斜的程度。同时,在 Celeborn 侧把 celeborn.worker.partitionSorter.batchSize 调小,让单个 Worker 上的 Partition 数量更分散,把热点数据的写入压力尽量打散。经过这两轮调整,热点问题缓解了不少,但数据倾斜特别严重的作业,仍然需要业务 SQL 层面改写。
6.2 Push 超时排查思路
Push 超时是引入 Celeborn 后最常见的报错之一。我们的排查路径是这样的:
第一步看网络,确认 Map 端所在节点到 Worker 节点的网络延迟和丢包率是否正常。第二步看 Worker 负载,如果 Worker 的 CPU 或磁盘 IO 已经打满,Push 请求排队是必然的。第三步看参数配置,是不是客户端设置的重试次数太少、超时时间太短。第四步看版本一致性,把客户端和服务端的版本都核对一遍。
大多数情况下,Push 超时都是前三类原因,版本不一致的问题在统一 jar 包版本之后已经很少出现了。这里建议大家在集群加机器的初期,一定要先做一次打通测试,再开放给业务使用,否则 Shuffle 数据和业务流量互相影响,定位问题会很头痛。
6.3 副本数带来的磁盘翻倍问题
PB 级作业的数据体量本身就很夸张,再加上 2 副本,磁盘空间的占用直接翻倍。如果 Worker 节点的磁盘容量预留不足,很容易出现磁盘满的尴尬情况。
我们的应对措施是给 Worker 节点挂载尽可能多的 NVMe SSD,并把磁盘使用率的告警阈值调到 75%。同时,把 Celeborn 的自动清理周期调短,确保作业结束后能及时释放磁盘空间。另外,建议对不同类型的作业设置不同的 celeborn.push.maxReplicate 值——核心作业设 2 副本,跑批中允许重算的作业可以降为 1 副本,达到空间和可靠性的动态平衡。
6.4 小作业开销变大的情况
这个现象很有意思。原本小作业走本地 Shuffle 只需要几秒钟,切换 Celeborn 后反而要多花十几秒。原因很简单:小作业的数据量不足以抵消推送和拉取的网络开销,Celeborn 的优势体现不出来。
我们的处理方式是配置化控制,让只有数据量超过一定阈值的作业才走 Celeborn。具体做法是在 Spark 的启动脚本里通过参数判断,如果作业预估输入数据量小于 50GB,就使用原来的 Shuffle Manager。这个阈值是根据我们集群的实际情况定的,每个团队最好结合自己的网络和存储条件调整。
7. 一点个人体会
回头看这段 PB 级 Shuffle 优化的过程,我最大的感受是:没有银弹,只有取舍。Celeborn 不是万能的,它引入了额外的部署组件和运维成本,也并非在所有场景下都能带来正向收益。但对于我们这种动辄跑几百 TB、上 PB 数据作业的平台来说,它解决的核心问题——小文件压力、故障重算代价、Shuffle 集群级性能瓶颈——是传统方案根本无法回避的。
如果你正打算在自己的集群上试 Celeborn,我建议先在测试环境完整跑通一遍 Spark 和 Flink 的集成流程,配好告警和监控,再逐步放量。另外,千万不要跳过小文件合并和压缩算法的验证,这两个环节往往是收益最明显的地方。踩过几次坑之后你会发现,Shuffle 优化说到底不是某一个参数的魔法,而是一整套架构设计、参数调优和运维保障体系的综合结果。
