ARTICLE DETAIL

建站实战干货

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

Apache Celeborn 0.5.x远程Shuffle服务实战:Spark/Flink集成与调优

2026/9/19 2:28:50 拓冰建站 浏览量
Apache Celeborn 0.5.x远程Shuffle服务实战:Spark/Flink集成与调优 先说一个判断如果你的 Spark/Flink 集群正在被 shuffle 拖慢或者你在为节点本地磁盘故障头疼Apache Celeborn 0.5.x 可能是你这两年最值得花时间研究的基础设施之一。Apache Celeborn 是一个远程 Shuffle ServiceRSS它在计算框架和存储之间插入了一层专门的 shuffle 数据服务把原本落在 Executor 本地磁盘的中间结果统一收走。0.5.x 是整个项目在 Apache 孵化器阶段逐渐稳定下来的一条版本线从部署方式、客户端 API 到生态兼容都形成了比较统一的规范。我在这条版本线上从零搭过环境、对接过 Spark 3.5 和 Flink 1.17、也压过线上任务这篇学习文档就把这段时间踩过的点和验证过的方案一次讲清楚。如果你正在做技术选型、刚开始接触远程 shuffle 概念或者已经部署了旧版本但想了解 0.5.x 的差别这篇内容都会帮得上忙。内容偏实操但原理部分我也会交代清楚尽量让“知其然”和“知其所以然”同时发生。1. 为什么需要 Celeborn传统 Shuffle 的痛点与远程 Shuffle 的解题思路1.1 传统 Shuffle 到底慢在哪先不急着解释 Celeborn 的组件得先搞清楚它解决的问题否则后面所有配置你都会觉得莫名其妙。Spark 和 Flink 这类分布式计算框架在算子之间传输数据时都有 shuffle 这个环节。以 Spark 为例Shuffle Write 阶段每个 map task 会把结果写到执行节点本地磁盘上的临时文件到了 Shuffle Read 阶段reduce task 再通过网络去各个节点拉取属于自己的那部分数据。这个流程本身没有大问题问题出在规模上来之后会放大几个毛病。一是本地磁盘 IO 和容量压力非常大。一个几十 TB 的批处理作业跑完shuffle 中间数据可能占掉每个节点几百 GB 甚至上 TB 的空间。你要提前给每台机器预留这么多磁盘而且读写都在同一批磁盘上IO 竞争非常严重。二是数据倾斜和故障恢复放大资源浪费。如果某个 map task 失败了Spark 默认要重新计算这个 task 对应的整个 shuffle 文件计算成本很高。三是集群调度和资源利用率被绑死。shuffle 数据在哪个节点下游任务如果跨节点拉取就会产生大量网络传输如果上游节点宕机就涉及重新计算这在云原生和弹性伸缩场景里特别难受。还有一个经常被忽略的问题动态资源分配。Spark 开启 Dynamic Allocation 之后executor 可能会中途释放或新增。传统 shuffle 中shuffle 数据写在本机executor 退出之前必须把 shuffle 文件处理好否则下游没法读。这就导致“想缩容但不敢缩”“想快速处理故障但 shuffle 在拖后腿”的尴尬局面。1.2 远程 Shuffle 的核心思路远程 Shuffle 的思路非常直接把 shuffle 数据的存储和计算节点解耦。引入一组独立的 Shuffle Service 节点map 阶段写数据不再落本地磁盘而是通过网络推送到 Shuffle Servicereduce 阶段直接从 Shuffle Service 拉数据。这样做了之后第一个好处就是计算节点的本地磁盘压力大幅下降。你不需要再为 shuffle 预留大量本地磁盘空间executor 的磁盘 IO 也主要花在处理应用数据上。第二个好处是数据可靠性集中管理。服务节点可以配置多副本某台节点挂掉时 shuffle 数据仍然可用不需要像传统方式那样重新跑上游 task。第三个好处是集群弹性更容易实现。executor 释放与否不再和 shuffle 数据生命周期强绑定动态资源可以更激进。这个思路不是 Celeborn 独有的像 Facebook 的 Cosco、腾讯的 Firestorm、LinkedIn 的 Magnet 等都做过类似尝试。Celeborn 更早脱胎于阿里内部实践所以工程成熟度相对高开源后社区迭代也比较快。1.3 Apache Celeborn 0.5.x 在版本序列里的位置Celeborn 从很早期就叫 Remote Shuffle Service后来进入 Apache 孵化器改名 Celeborn。0.5.x 这条版本线我在学习过程中把它理解为“从可用到易用”的转折点。0.2.x、0.3.x 的阶段项目还在快速迭代很多配置项会变客户端 jar 的兼容矩阵也没有现在清晰。到 0.5.x主版本固定之后对应的 Spark 版本覆盖到 3.1/3.2/3.3/3.4/3.5Flink 覆盖 1.14 的多个小版本MapReduce 也有对应的支持。部署端的 Master/Worker 配置也收敛了运维文档补得比较全。0.5.0 作为这条分支的首个版本发布后0.5.1、0.5.2、0.5.3 主要在做稳定性和 bug 修复所以学习资料按 0.5.x 整体看是合理的。如果你是从旧版升级上来或者一步到位选型直接以 0.5.x 为学习基线非常合适。配置规范延续到后续版本学完不会白学。2. 0.5.x 架构拆解Master、Worker 与客户端如何协同2.1 三角色模型Celeborn 0.5.x 的运行时架构分三块理解起来比大多数分布式系统简单。Master 是整个集群的控制面。它负责管理 worker 节点的心跳、资源状态、shuffle 分区分配和元数据。生产环境通常部署 3 个 Master通过 Raft 协议选主和复制元数据。客户端找 Master 获取 worker 列表、注册 shuffle、提交分区信息这些控制信令都走 Master。Worker 是数据面。它接收上游推过来的 shuffle 数据写入本地存储同时响应下游的 fetch 请求。Worker 可以同时跑很多个每台机器上可以部署一个或多个实例存储目录可以在配置里指定多个。生产里通常把 Worker 部署在独立的机器集群上和计算节点分开也可以混部看你的资源情况。客户端不是一个独立进程而是嵌入在 Spark Executor、Driver 或者 Flink TaskManager 里的组件。它负责把 map task 产生的数据按分区切好批量推送到对应 Worker也负责从 Worker 拉数据给 reduce 端。客户端还包含一个 lifecycle manager跟 Master 通信完成 shuffle 注册和释放的协调。这三个角色加起来本质上就是把原来“每个 Executor 自己管 shuffle 文件”变成了“Master 管元数据Worker 管存储客户端管传输”的三方协作模型。2.2 Shuffle 数据通路与生命周期在一个 Spark 任务里Celeborn 介入 shuffle 之后数据流是这样的。第一步Driver 端的 ShuffleManager 初始化时通过 lifecycle manager 向 Master 注册 application申请 shuffle id。Master 根据集群当前 worker 资源返回可用的 worker 列表。第二步每个 map task 写完自己的输出后不是写到本地 disk而是通过客户端把数据按 reduce 分区切好推送到分配好的 Worker 上。Worker 把数据写入存储目录先写内存缓冲区再异步刷入底层文件。第三步reduce task 开始读数据时客户端拿着 shuffle id 和 reduce 分区号去 Master 查对应的 Worker 位置然后直接向 Worker 发起 fetch 请求拿到数据。整个过程里有一个关键设计就是分区到文件的切分粒度。Celeborn 在 Worker 端把每个 shuffle 的每个 reduce 分区数据落成独立的文件段这样 fetch 的时候可以精准定位不用全量扫描。如果配置了多副本Worker 会把同一份数据复制到另一个 Worker 上fetch 端可以做选择或 failover。生命周期管理上Driver 结束时会通知 Master 释放该 application 占用的 shuffle 空间。Master 再通知 Worker 做文件清理。如果 Driver 异常退出Master 有超时机制超过阈值后自动回收资源避免数据泄漏撑爆磁盘。2.3 存储模型与高可用设计0.5.x 的存储层支持本地文件系统和 HDFS 两种底层存储这是一个非常实用的设计。本地文件系统模式下Worker 直接把数据写到配置的磁盘目录读写性能最好适合对延迟敏感和纯离线批处理场景。HDFS 模式的逻辑是把 shuffle 数据落到底层 HDFS好处是存储容量可以做的非常大而且天然有 HDFS 的副本机制适合数据量巨大但对实时性要求不苛刻的场景。两种模式不是互斥的你在 worker 存储目录里配置不同前缀或者不同实现即可。可靠性和高可用方面Celeborn 0.5.x 做了两层保障。Master 侧靠 Raft 多副本选主客户端会自动重连新的 Leaderworker 注册和元数据查询都不需要手动切换。数据侧靠多副本推送默认可以配置副本数为 1 或 2。副本数为 2 的时候客户端会等两个 Worker 都返回成功才认为写入完成读端可以做均衡选择。这样单台 Worker 挂掉就不太会影响作业shuffle 数据依然能从另一副本读取。我个人的理解是这套架构跟计算框架本身不耦合业务逻辑所以可扩展性很不错。Master 和 Worker 之间通过 gRPC 通信客户端到 Worker 之间走 Netty 长连接整体性能和稳定性在 0.5.x 已经比较成熟。3. 本地环境快速搭建 0.5.x 集群3.1 下载与包结构学习 Celeborn 0.5.x 最好的方式是先在本地或者单台上把整套服务跑起来。不需要很大的集群一台 4 核 8G 的虚机就够做功能验证。去 Apache 官网下载二进制包解压之后你会看到几个重要目录bin/放启动停止脚本conf/放配置模板spark/和flink/目录下放各个版本的客户端 jarsbin/放 master/worker 管理脚本logs/是运行日志目录。下载时建议认准 Apache 官方 release 或者国内镜像站的稳定版本不要从不明渠道拿编译包。0.5.x 的二进制包命名大概是apache-celeborn-0.5.x-bin.tgz里面已经包含 master、worker 的启动脚本以及常用客户端 jar。3.2 最小配置启动 Master 和 Worker解压完成后先把配置模板复制出来cd apache-celeborn-0.5.x-bin cp conf/celeborn-defaults.conf.template conf/celeborn-defaults.conf最小化配置里我建议至少设置这几个参数celeborn.master.endpoints localhost:9097 celeborn.worker.storage.dirs /data/celeborn/worker1 celeborn.worker.flush.buffer.size 256kceleborn.master.endpoints是客户端和 worker 找 Master 的入口地址本地单机填 localhost 就行多个 Master 用逗号分隔。celeborn.worker.storage.dirs是 Worker 写数据的目录最好用独立挂载点别和系统盘混在一起。celeborn.worker.flush.buffer.size控制刷盘缓冲大小测试环境小一点没关系生产再调。启动 Mastersbin/start-master.sh启动 Workersbin/start-worker.sh启动完可以看日志确认状态。Master 默认会起两个端口一个 RPC 通信端口默认 9097一个 web UI 端口默认 9098。浏览器打开http://localhost:9098如果能看到集群信息和注册的 worker 列表说明 Master 起来了。Worker 没有独立 web 页面但会在日志里打印注册结果。看到类似Registered successfully的信息就说明 Worker 连上 Master 了。3.3 验证集群是否可用很多人启动完服务就急着去接 Spark我建议先做一次简单验证花两分钟能排除很多后续问题。可以先看 web UI 上的 worker 列表确认 worker 状态是online。然后看日志文件logs/celeborn-worker.out和logs/celeborn-master.out确认没有报错堆栈。再检查存储目录是否被创建出来权限是否正常。如果存储目录创建失败后面所有推送都会失败。一个比较容易被忽略的点是端口连通性。如果你的测试环境有防火墙或者你后面要让远程的 Spark 集群连这个 Celeborn 集群一定要确认 9097 和 8090 这类端口对外可达。客户端连接走的是 RPC 端口数据推送走的是 Netty 端口两个不通都不行。检查端口可以用telnet localhost 9097也可以看 worker 日志里绑定的实际端口。4. Spark 与 Flink 集成完整实操4.1 Spark 集成与客户端 Jar 选择Celeborn 0.5.x 对 Spark 的支持版本已经非常清晰每个 Spark 小版本都有对应的客户端 jar。以 0.5.3 为例spark/目录下能看到celeborn-client-spark-3.1.jar、celeborn-client-spark-3.2.jar、celeborn-client-spark-3.3.jar、celeborn-client-spark-3.4.jar、celeborn-client-spark-3.5.jar这类文件。选 jar 版本时有一个原则按 Spark 主版本和次版本严格匹配不要偷懒用低版本的 jar 去跑高版本的 Spark。虽然 Spark 的 API 通常向后兼容但 shuffle manager 这类接口变化频繁对应关系错位会导致启动阶段就报ClassNotFoundException或者NoSuchMethodError。我以 Spark 3.5 为例展示启动时需要加的配置。在spark-defaults.conf里加这些内容spark.shuffle.manager org.apache.spark.shuffle.celeborn.SparkShuffleManager spark.serializer org.apache.spark.serializer.KryoSerializer spark.celeborn.master.endpoints localhost:9097 spark.shuffle.service.enabled false spark.dynamicAllocation.enabled falsespark.shuffle.manager是核心入口告诉 Spark 不用默认的排序 shuffle 管理器而改用 Celeborn 的实现。spark.serializer必须设成 Kryo因为 Celeborn 客户端序列化依赖 Kryo你用 Java serializer 在后面压缩和缓存环节容易出问题。spark.celeborn.master.endpoints是客户端找 Master 的地址和 Celeborn 服务端配置保持一致。spark.shuffle.service.enabled和spark.dynamicAllocation.enabled在生产里通常要关掉因为远程 shuffle 自己管理数据不需要 Spark 的 external shuffle service 和基于 shuffle 文件的动态分配机制。如果你用spark-submit跑作业还要把客户端 jar 加到--jars里或者提前放到 Spark 的jars/目录。我一般建议用--jars这样不影响其他作业升级也方便。4.2 Flink 集成要点Flink 的集成思路和 Spark 类似但没有 Spark 那么“一键切换”需要手动指定 shuffle 服务。0.5.x 对 Flink 的兼容覆盖了 1.14 到 1.19 左右的常见版本。使用方式是把对应的celeborn-client-flink-xxx.jar放到 Flink 的lib/目录然后在flink-conf.yaml里设置 shuffle service 相关参数。核心配置大概是shuffle-service-factory.class: org.apache.flink.runtime.io.network.celeborn.CelebornShuffleServiceFactory同时把 Celeborn 服务端端点配置传给 TaskManager一般是设置celeborn.master.endpoints。Flink 集群内所有 TaskManager 都要能从网络访问到 Celeborn 集群。Flink 集成有一个需要注意的地方Flink 的 shuffle 模型和 Spark 不完全一样Celeborn 对不同 Flink 小版本的适配方式也有差异。所以建议查看官方 release notes 里对应版本的支持说明不要只看一个博客就盲目上线。我自己在 Flink 1.17 上验证过基本流程是通的但 Flink 1.14 和 1.15 在某些参数名上有细微区别。4.3 关键参数与常见配置组合掌握了基础切换之后重点就是理解参数之间的配合。Celeborn 0.5.x 的客户端参数分成几组我挑最重要的几个说。推送相关的参数里spark.celeborn.client.push.buffer.size控制每个分区的推送缓冲大小默认值不算大生产环境如果网络好可以适当调大减少小包数量。spark.celeborn.client.push.merge.enabled合并小数据包对于高频小 shuffle 非常有帮助开启后网络 overhead 明显降低。拉取相关的参数里spark.celeborn.client.fetch.maxReqsInFlight控制 fetch 并发数太大会导致 Worker 过载太小会拉低吞吐。建议从默认值开始压测时观察服务端 CPU 和网络再逐步调。压缩是一个很重要的优化点。spark.celeborn.client.compression.codec可以选lz4或zstd。生产我验证下来 zstd 压缩比更高CPU 开销也可接受对磁盘和网络节省明显。但要注意设置 compression 之后客户端和服务端版本要兼容同一个 Celeborn 集群里如果混着不同客户端版本压缩配置最好统一。还有一个容易踩坑的地方是spark.celeborn.shuffle.chunk.size。它决定 fetch 时单次传输的数据块大小很多文档没有详细说。数据块太小RPC 次数变多太大内存占用升高。一般保持默认只有在 shuffle 数据集特别大并且网络带宽充足时才考虑调更大。5. 运维监控与问题排查实录5.1 服务健康检查与日志线索Celeborn 0.5.x 的运维第一步是熟悉日志。Master 的日志最值得关注的信息点有leader 切换事件、application 注册和释放记录、worker 心跳超时。如果看到频繁的 leader 切换需要检查 Master 节点之间的网络和时钟。如果看到 worker 掉线后重连要检查 worker 所在节点的负载和磁盘。Worker 的日志通常包含三类内容启动时注册信息、推送和 fetch 请求异常、存储目录不可写。一个高频问题就是磁盘写满导致 worker 把目录标记为不可用这种情况在日志里会有明确提示生产环境要第一时间处理。Master 的 web UI 在运维里很好用能直接看到当前集群有几个 worker、每台 worker 的状态、存储使用率、当前存活的 shuffle application。我建议把它加入你的监控面板至少每 5 分钟采集一次状态。5.2 常见故障场景与处理办法我在实际使用中遇到的故障整理成了一张速查表方便你们对照排查。现象可能原因处理办法客户端连不上 Mastermaster.endpoints 填错、端口不通检查配置和防火墙telnet 测试端口确认 Master 进程存活Worker 启动后反复掉线磁盘目录不可写、worker 心跳超时检查存储目录权限和空间调整 heartbeat 时间参数作业提交后 ClassNotFound客户端 jar 与 Spark 版本不匹配严格按照 Spark 次版本选择 client jar重新提交作业Shuffle 写入阶段超时worker 负载高、网络抖动、推送缓冲区过大观察 worker CPU/网络适当调小推送缓冲检查是否单 worker 热点Fetch 时数据缺失副本数配置为 1且对应 worker 宕机设置celeborn.worker.replicate为 true或在关键作业上保证多副本磁盘空间暴增历史 application 未及时清理、清理超时时间过长检查 master 的清理阈值手动调用清理接口或重启 workerKryo 序列化报错用户自定义类没有注册 Kryo为自定义类型注册 Kryo 类或关闭 Kryo 相关优化从业务侧排查这里说一个我踩过的坑celeborn.worker.monitor.disk.enabled这个开关。默认是开启的Worker 会定期检查磁盘健康如果某个目录写不进去会自动把该目录置为不可用。这个机制本身是保护但如果你配置了多个存储目录其中一个目录是挂载不稳定的网络盘会导致 Worker 频繁切换可用目录性能波动很大。生产环境尽量用稳定的本地盘并且给磁盘监控设置合理的阈值。另一个坑跟 Spark 的 external shuffle service 有关。很多现有 Spark 集群开了spark.shuffle.service切换 Celeborn 之后如果忘记关闭会出现新旧两个 shuffle 通道打架的问题现象是部分数据能从 Celeborn 读部分数据走本地最终一致性出问题。所以切换 Celeborn 时一定要把 external shuffle service 关掉并且确保spark.dynamicAllocation.enabled也是关闭状态或者经过完整验证再开。5.3 版本升级与滚动重启0.5.x 内部的小版本升级比如从 0.5.0 升到 0.5.2通常会包含 bug 修复和安全性更新建议跟上。升级节奏我建议这样先升级 client jar让线上作业跑在旧集群上验证新客户端兼容性然后逐个升级 Worker滚动重启最后升级 Master。因为 Master 有 Raft 选主滚动重启时会自动切换 Leader不会对运行中的作业造成明显影响。升级之前一定要看目标版本的 release notes重点看有没有配置项改名、默认值变更和不兼容项。升级过程中有一个很重要的注意点客户端和服务端不要跨大版本混跑。比如客户端用 0.5.x服务端升级到 0.6.x这个组合大多数情况下没问题但社区不一定保证兼容尤其是协议字段有变化的时候。能对齐版本就对齐版本别图省事。6. 性能调优与实践心得6.1 资源规划与服务端调优生产部署 Celeborn 0.5.x资源规划不能拍脑袋。我一般按三个维度估算磁盘容量、内存、网络。磁盘容量方面先粗略估算一下集群每天跑批的 shuffle 数据总量。假设一天 shuffle 总量是 10TB保留两天的数据滚动清理那么 Celeborn 集群至少要 20TB 可用存储。如果开启多副本再乘以副本数。这里的安全系数建议留到 1.5 倍以上因为高峰期数据量往往比均值高很多。内存方面Worker 进程除了堆内存还有大量堆外内存用于网络缓冲。单台 Worker 如果承载 20 个并发作业堆内存建议 16GB 以上堆外内存按每个客户端连接约几 MB 计算。不要在一台超大内存机器上部署太多 Worker 实例单机多实例反而会增加整体管理复杂度我建议一台机器一个 Worker 为主。网络方面千兆网卡在数据量大的时候会非常吃力生产环境强烈建议万兆网络。如果条件不允许就要在客户端侧调小推送并发和缓冲区避免把网络打爆导致超时重试。服务端还有几个参数值得调。celeborn.worker.fetch.chunk.size决定 fetch 时返回的数据块大小磁盘较弱时调小可以减少单次读取压力网络较弱时调大可以减少 RPC 次数。celeborn.master.ha.raft.node相关参数只在多 Master 模式下有效单机模式不用管。6.2 客户端参数组合的实测经验调优不能只调服务端客户端参数的组合对最终效果影响更大。我拿 Spark 3.5 跑 TPC-DS 的部分 SQL 做过对比。在默认参数下Celeborn 模式比本地 shuffle 模式网络传输量增加因为数据从计算节点多走了一跳但总耗时不一定劣化因为避免了本地磁盘的随机 IO。当我把spark.celeborn.client.push.merge.enabled打开并且把spark.celeborn.client.push.buffer.size调大到 256k 时整体耗时比默认参数提升了 10% 到 20%因为网络包数量锐减服务端 CPU 也降了。压缩参数方面zstd 在数据倾斜明显的场景下收益更大。倾斜场景里大量数据集中在少数 key 上压缩率能到 4:1 甚至更高不仅能减少网络传输量还能减少 Worker 端存储占用。我在验证中zstd 模式下 CPU 平均占用比 lz4 高 5% 左右但 shuffle 总量下降 30% 以上整体划算。还有一个容易被忽视的参数是spark.celeborn.client.push.timeout。网络环境不太稳定的集群超时时间设太短会导致大量重试反而把服务端打挂。我建议默认值先跑观察日志里 push timeout 出现频率如果频繁出现再逐步调大。6.3 从传统 Shuffle 迁移的平滑策略最后聊一下迁移策略这是很多团队最犯难的地方。我推荐的路线是先做影子验证再做小流量切换最后全面铺开。影子验证的意思是保持现有集群跑传统 shuffle另外搭一个 Celeborn 测试集群用备份数据或者抽样数据跑相同作业对比稳定性和性能。这个小集群不用太大能支撑几个作业并发就可以但它能帮你提前发现 jar 版本、参数配置、权限问题。小流量切换阶段挑两三个非核心作业用单独跑的 Spark 集群或者提交参数切换成 Celeborn 模式观察一两天。这里特别要注意业务侧对 shuffle 性能不敏感。如果作业本来就跑得很快切换 Celeborn 带来的提升不明显反而增加运维面性价比不高。全面铺开时我建议按队列或者业务域分批次切而不是所有作业一天全切。因为你不知道哪个作业会触发 Celebdoorn 的新问题一旦出问题受影响面越小越好。把 Celeborn 集群的监控和大盘提前做好出现异常能快速定位是哪个 worker、哪个 application、哪类请求。迁移后还有一件事别忘关闭旧的外部 shuffle service 依赖。很多现有 Spark 脚本里会显式配置spark.shuffle.service.enabledtrue切换 Celeborn 时需要同步清理。如果线上线下参数管理不统一很容易出现“主脚本已经切了 Celeborn某个调度模板还带着旧参数”的情况。6.4 我对 0.5.x 的整体评价最后说点主观判断。Celborn 0.5.x 不是一个“装了就立刻起飞”的组件它解决的是规模和可靠性问题。如果你的集群节点少、shuffle 量不大、本地磁盘很健康传统 shuffle 完全够用没必要引入新组件。但如果集群规模到几十台以上shuffle 导致磁盘和网络成为瓶颈或者你计划上动态资源、弹性节点这类特性Celeborn 的收益会非常明显。0.5.x 最大的优势在于兼容和稳定。我在学习和部署过程中问题大多出在版本匹配和参数理解上而不是组件本身的 bug。在社区活跃度上0.5.x 的文档也不像早期版本那么难啃基本配置照着就能跑通。如果现在有 0.6 或更晚的版本出来我的建议是仍然先以 0.5.x 为基线做学习和功能验证因为它是目前生态文档最齐、踩坑分享最多的一条主线。等新版本验证充分了再评估迁移也不迟。