十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Uniffle:统一Shuffle引擎如何解决Spark大数据作业的稳定性难题

Apache Uniffle:统一Shuffle引擎如何解决Spark大数据作业的稳定性难题 1. 为什么大数据工程师都在聊 Uniffle这两年大数据圈子里统一 Shuffle 引擎这个词出现得越来越频繁而每次聊到这个话题Apache Uniffle 基本是绕不开的那个名字。如果你正在维护一套 Spark 或者 MapReduce 生产集群或者每天被 NodeManager 上的 Shuffle 临时文件搞到磁盘告警那这篇文章值得你花几分钟看完。先说清楚 Uniffle 到底是什么。Apache Uniffle 是一个开源的、统一的 Shuffle 引擎项目前身是腾讯在 2020 年开源的 RSSRemote Shuffle Service后来在 2021 年捐献给 Apache 基金会孵化2023 年正式成为 Apache 顶级项目。它的定位很直接把原本跑在计算节点本地、跟 YARN NodeManager 耦合在一起的 Shuffle 过程拆出来放到独立的集群中执行。传统 Shuffle 的问题做过大数据的人应该都有感触。以 Spark 为例每个 Executor 算完 Map 任务后要把中间结果先写到本地磁盘然后 Reduce 任务再从各个节点把属于自己的那份数据拉过来。这个过程中会产生海量的小文件一个跑 5000 个 Map 任务的 Job如果有 100 个 Reduce那就要产生几十万个临时文件。文件一多NodeManager 的文件句柄数飙升磁盘 IO 被打满GC 压力变大整个集群的稳定性都会被拖下水。Shuffle 一旦出现故障问题就更大。某个 NodeManager 因为 Shuffle 文件写太多挂掉那它上面所有未完成的 Shuffle 数据就全丢了所有依赖这批数据的 Reduce 任务都要重新去拉去等整个作业直接卡死。有些团队为了规避这种问题只能把 Shuffle 相关的临时目录分散到多块盘上但盘再多也撑不住大作业的临时文件洪峰。Uniffle 解决的就是这一整套问题。它把 Shuffle 数据的写入和读取从计算节点本地转移到了独立的 Shuffle Server 集群规格可以单独规划磁盘可以单独管理IO 压力不再跟计算资源抢。Coordinator 负责管理和调度Shuffle Server 负责存储和分发Map 任务直接把数据推到 Shuffle ServerReduce 任务从 Shuffle Server 拉取计算节点本地不再保存任何 Shuffle 中间数据。这篇文章适合谁看如果你是平台工程师、大数据架构师或者正在被 Spark 作业的 Shuffle 稳定性问题折磨的开发者那 Uniffle 的架构思路和接入方式值得你认真研究一遍。就算你现在只是用着开箱即用的 EMR 或者 CDH了解清楚它的核心原理对后续排查 Shuffle 相关性能问题也会有帮助。2. 核心架构设计Coordinator 与 Shuffle Server 的分工逻辑要理解 Uniffle 为什么能解决传统 Shuffle 的问题得先把它的架构拆开看。Uniffle 整体上只有两个核心组件Coordinator 和 Shuffle Server。我刚开始接触的时候也觉得这个设计简单得有点意外但实际用下来会发现这种少即是多的架构恰好是它稳定和易维护的基础。2.1 Coordinator集群的调度中枢Coordinator 在 Uniffle 里的角色类似 YARN 里的 ResourceManager或者 Kafka 里的 Zookeeper。它本身不存储 Shuffle 数据主要负责两件事管理 Shuffle Server 的注册与状态以及给客户端包括 Map 端和 Reduce 端分配可用的 Shuffle Server 列表。Coordinator 会定期收集所有 Shuffle Server 的心跳信息包括当前负载、磁盘剩余空间、连接数等。当客户端请求 Shuffle Server 列表时Coordinator 会根据配置的分配策略返回一批合适的节点。默认的分配策略相对简单主要是看磁盘使用率和连接数但实际生产中我一般会自己扩展一下分配策略把机房或者机架信息考虑进去避免 Shuffle 数据全部打在同一批机器上。因为是集中式调度Coordinator 自身必须高可用。Uniffle 支持部署多个 Coordinator 实例组成集群它们之间通过 ZooKeeper 或者内置的 Raft 协议做选主和元数据同步。不过说实话Coordinator 的元数据非常轻主要就是 Server 列表和心跳状态所以瓶颈基本不在它的存储上而在于连接数的处理能力。遇到大集群或者高峰期建议给 Coordinator 单独预留足够的堆内存并且把 GC 参数调好。2.2 Shuffle Server数据读写的主战场Shuffle Server 是 Uniffle 真正干重活的地方。每个 Shuffle Server 负责接收来自 Map 端的数据写入并存储到本地然后等待 Reduce 端来拉取。从职责上看它就像一个中转仓库Map 端把货送过来Reduce 端来提货中间不需要计算节点存任何东西。Shuffle Server 内部使用 Netty 处理网络请求数据会被追加写入到本地预分配的大文件中而不是像原生 Shuffle 那样每个 Map 任务每个 Reduce 分区都生成一个单独的小文件。这一下就把文件数量从百万级降到了个位数级别文件句柄和磁盘 IO 的压力自然就下来了。每个 Shuffle Server 还可以挂载多块磁盘Uniffle 会以目录为单位做数据均衡轮询写入不同的目录。跟 HDFS 的 DataNode 比Shuffle Server 的读写路径简单很多不需要副本复制、不需要 NameNode 参与纯本地磁盘操作所以 IO 效率很高。还有一个值得注意的设计Shuffle Server 是有状态的它保存着 Shuffle 数据如果它挂了上面的数据会丢失。这与 HDFS 的高可靠设计完全不同因此 Uniffle 在故障恢复上要依赖上层的重试机制。为了减少这种风险生产环境通常会部署多台 Shuffle Server并且利用 Coordinator 的分配策略把单个作业的数据尽量打散到多台机器上。2.3 为什么这个设计比传统方案更可靠传统 Shuffle 的最大痛点是计算和存储耦合。Map 任务写数据到本地磁盘Reduce 任务再从各个节点远程拉取数据这本质上是 N 对 N 的拉取模型。一旦某个节点上的 Shuffle 文件损坏或丢失所有等待该数据的 Reduce 任务都会失败。Uniffle 把 Shuffle 数据集中到 Shuffle Server 后数据拉取模型变成了 N 对 MM 是 Shuffle Server 的数量远小于 N。计算节点不再参与 Shuffle 数据的保存节点宕机不会再导致 Shuffle 数据丢失。而对于 Shuffle Server 的故障因为作业的数据分散在多台 Server 上单台故障只影响一部分数据Map 端可以基于重试机制或上游的重新计算来恢复整体影响面小很多。另外集中式 Shuffle 还带来了一个隐藏收益集群的弹性伸缩变得更容易了。计算节点不需要预留大块本地磁盘给 Shuffle可以随时弹性扩缩容因为 Shuffle 数据不落在计算节点上缩容时不用担心数据丢失。这一点对于云原生场景特别有价值。3. 深入理解 Uniffle 的 Shuffle 工作流程架构搞清楚之后接着看数据流转的完整过程。Uniffle 对上层应用暴露的接口兼容 Hadoop MapReduce 和 Spark 的 Shuffle 协议所以接入时不需要修改计算框架的内部逻辑而是通过插件方式替换掉原来的 Shuffle Manager。下面以 Spark 为例把一条 Shuffle 数据从 Map 端到 Reduce 端的完整路径走一遍。3.1 作业启动Shuffle Server 的分配与注册Spark 作业提交后在 Driver 端会初始化 Uniffle 的 ShuffleManager。此时 Driver 会向 Coordinator 请求分配一批 Shuffle Server。Coordinator 接收到请求后根据当前所有 Shuffle Server 的负载情况返回一个可用的 Server 列表。然后在 Executor 启动时Map 端的 Executor 也会从 Coordinator 获取同一份 Server 列表并对每一个 Shuffle Server 建立一个长连接。后续的数据写入都会复用这个连接避免频繁建连的开销。如果某个 Server 在作业运行期间挂掉Executor 会重新向 Coordinator 请求新的 Server 列表并更新本地缓存。这里有一个常见误区很多人以为 Uniffle 的 Server 列表是每个作业单独分配的其实它是按照Shuffle ID 维度来分配的。同一个作业的不同 Shuffle Stage可能会分配到不同的 Server 集合这样可以让负载更均衡。3.2 Map 端写入从本地文件到远程推送在原生 Spark Shuffle 中Map 任务会把输出数据写到本地磁盘然后报告给 Driver 一个 MapStatus里面记录着每个 Reduce 分区对应的文件位置和大小。而使用 Uniffle 后Map 任务的输出不会写本地磁盘而是直接写入内存缓冲区达到阈值后通过 Netty 批量推送给对应的 Shuffle Server。Map 端数据先经过 Partitioner 分区每个分区在内存中积攒一个批次这个批次会作为一次网络请求发送给目标 Shuffle Server。Shuffle Server 收到数据后会先写入内存再异步刷到本地大文件中。这样做的好处是网络请求次数大幅减少磁盘随机写变成了顺序追加写。这个阶段最怕的是数据倾斜某些分区数据量特别大内存缓冲区会先满导致频繁刷写。我建议在接入 Uniffle 时把 Spark 的 Shuffle 相关内存参数调大一点或者适当增加分桶数量让数据分布更均匀。3.3 Reduce 端拉取并发读与数据校验Reduce 任务启动后会从 Driver 获取 MapStatus不过这个 MapStatus 记录的已经是 Shuffle Server 上的位置信息。Reduce 任务会向 Coordinator 确认这些 Server 仍然可用然后并发地从多台 Shuffle Server 拉取属于自己的分区数据。Uniffle 的 Reduce 端拉取做了很多优化比如读缓存、合并请求、校验和过滤等。数据从 Shuffle Server 读取时会先检查本地索引文件确定数据块在数据文件中的偏移量和长度然后直接从对应位置读取目标数据不用全文件扫描。数据拉到 Reduce 端之后Uniffle 会做完整性校验确保网络传输中没有丢包或者错位。校验失败的数据块会自动重新拉取这个重试逻辑在 Reduce 端是透明的不需要用户介入。不过需要注意Reduce 端的并发度不能设置得过于激进否则大量并发连接同时打向 Shuffle Server网络和内存都可能成为瓶颈。3.4 文件布局大文件加索引的设计思路说到存储设计这是 Uniffle 让我觉得最妙的地方。原生 Shuffle 会为每个 Map 任务和每个 Reduce 分区生成一个文件片段小文件数量是 Map 数乘以 Reduce 数。Uniffle 则让每个 Shuffle Server 为每个 Shuffle 任务只生成一个数据文件和一个索引文件。数据文件负责存储实际的数据内容按分区-批次的顺序追加写入。索引文件则记录每个分区下每个批次的数据在数据文件中的偏移量和长度。Reduce 端拉取数据时先读索引文件找到目标偏移范围再定位数据随机读效率非常高。这种一个大文件加一个索引文件的设计不仅是减少文件数量那么简单。它让磁盘空间的管理变得容易了数据文件可以预分配大小文件系统碎片减少而且清洗过期数据的时候直接删掉整个目录就行清理效率远超几百万个小文件的逐条删除。4. 接入 Uniffle 的完整实操记录理论部分讲完就该动手了。我以一个 Hadoop 3.3.4 加 Spark 3.3.1 的环境为例把 Uniffle 从部署到接入的完整过程捋一遍。过程中我把实际踩过的坑也一并写出来帮你少走弯路。4.1 前提条件与版本选择在开始之前先确认你的环境满足这些条件JDK 8 或 JDK 11Hadoop 2.10.x 或 3.xSpark 2.4.x 或 3.x如果你要用 Spark 的话ZooKeeper用于 Coordinator 高可用也可以使用内置 Raft关于版本选择我的建议是尽量用最新稳定版因为 Uniffle 的迭代速度比较快早期版本有一些性能和稳定性问题到 0.8.0 之后的版本在生产环境才表现得比较成熟。我用的是 0.8.0整体体验不错。4.2 编译与部署 Shuffle ServerUniffle 的 release 包可以直接在 Apache 官网下载如果你想修改配置或者自定义某些功能也可以从 GitHub 拉源码自己编译。编译命令很简单git clone https://github.com/apache/incubator-uniffle.git cd incubator-uniffle ./build.sh --skipTests编译完成后在dist目录下会有完整的部署包。把部署包分发到你的 Shuffle Server 机器上然后修改conf/rss-env.sh主要配置以下内容# Java 环境变量 JAVA_HOME/usr/local/jdk1.8.0_202 # Shuffle Server 的 JVM 参数堆内存建议根据单机 Shuffle 数据量来定 RSS_SERVER_JVM_OPTS-Xms32g -Xmx32g -XX:UseG1GC然后修改conf/server.confrss.rpc.server.port19999 rss.rpc.server.http.port19998 rss.storage.typeMEMORY_LOCALFILE rss.storage.basePath/data1/rss-data,/data2/rss-data,/data3/rss-data rss.server.memory.shuffle.highWaterMark.percentage60.0 rss.server.memory.shuffle.lowWaterMark.percentage40.0 rss.server.app.expired.withoutHeartbeat60000 rss.coordinator.quorumcoordinator1:19999,coordinator2:19999这里重点解释几个关键参数rss.storage.type存储类型可以选择MEMORY_LOCALFILE内存加本地文件或者LOCALFILE纯本地文件。我推荐前者因为热数据可以缓存到内存读性能更好。rss.storage.basePath指定存储目录可以配置多个路径Uniffle 会做目录级别的负载均衡。盘越多IO 并发度越高。水位参数控制内存使用的上下限避免 Server 内存溢出。启动 Shuffle Server 只需要执行bin/start-shuffle-server.sh启动后通过jps检查 RSSShuffleServer 进程是否存在再去看日志确认注册是否成功。4.3 部署 CoordinatorCoordinator 的部署更轻量。同样把部署包分发到 Coordinator 节点修改conf/coordinator.confrss.coordinator.server.port19999 rss.coordinator.http.server.port19998 rss.coordinator.select.strategyIO rss.coordinator.exclude.nodes.enabledtrue rss.coordinator.node.exclude.check.interval.ms60000关键参数中rss.coordinator.select.strategy决定了给客户端分配 Server 时的策略可选IO、MEMORY、CPU等。我一般用IO因为 Shuffle 场景下最稀缺的资源就是磁盘 IO。启动命令bin/start-coordinator.sh启动后用浏览器访问http://coordinator_ip:19998能看到当前注册的 Shuffle Server 列表和集群状态这个页面在排查问题的时候非常有用。4.4 配置 Spark 客户端接入客户端配置是接入过程中最容易出问题的地方。我当时的配置如下spark.jars /path/to/rss-client-spark3.jar spark.shuffle.manager org.apache.spark.shuffle.RssShuffleManager spark.rss.coordinator.quorum coordinator1:19999,coordinator2:19999 spark.rss.storage.type MEMORY_LOCALFILE spark.rss.client.read.buffer.size 16m spark.rss.client.write.buffer.size 64m spark.rss.client.send.size.limit 16m配置好之后提交一个简单的测试作业spark-submit --class com.example.ShuffleTest \ --master yarn \ --deploy-mode cluster \ --conf spark.shuffle.managerorg.apache.spark.shuffle.RssShuffleManager \ rss-demo.jar作业跑起来后通过 Coordinator 的页面观察 Shuffle Server 的读写吞吐。如果能看到数据持续流入说明 Map 端写入已经生效了。4.5 一个必调的客户端参数在接入后第一周我发现了一个问题当 Shuffle Server 负载过高时部分 Reduce 任务会出现长时间等待导致整个 Job 变慢。排查后发现是客户端重试的超时时间设置得太保守。spark.rss.client.read.retry.wait默认比较短在 Server 队列积压严重时容易触发超时。我把重试等待时间调到了 5 秒同时增加了重试次数到 10 次问题就缓解了。不过这里也要提醒一句参数不能无脑调大否则单次作业失败等待时间会很长。5. 常见问题排查与性能调优经验接入任何新组件都会遇到各种各样的问题。这里把我在生产环境使用 Uniffle 期间遇到过的典型问题整理成一个速查表每一个都是实际踩过坑之后总结出来的。5.1 系列故障排查清单现象可能原因排查方式解决方案Coordinator 页面无 Server 注册Server 启动失败或网络不通检查 Server 日志telnet 测试端口确认端口未占用检查防火墙策略作业长时间卡在 Shuffle 阶段数据倾斜或 Server 队列积压查看 Server 日志中队列长度检查 Map 端输出数据分布增加分区数调大高水位参数增加 Shuffle Server 节点Reduce 端大量重试网络抖动或 Server 文件损坏查看 Reduce 端错误日志调大重试次数检查磁盘健康状态作业报错找不到 Shuffle ServerCoordinator 分配策略过于集中查看 Coordinator 日志确认分配结果调大 Server 列表数量检查 Server 心跳是否正常磁盘 IO 仍然很高数据本地化率低或 Server 数量不足通过监控查看每台 Server 的 IO 指标扩展 Shuffle Server 节点增加挂载盘数量内存占用持续升高内存水位参数设置不合理观察 GC 日志查看堆内存使用曲线调低高水位、调高低水位启用 G1 回收器这里的核心排查思路是先看 Server 端日志再看 Coordinator 日志最后才看客户端日志。因为大多数问题的根因都出在服务端客户端只是暴露了症状。5.2 数据倾斜场景下的特殊处理数据倾斜在传统 Shuffle 中就是老大难问题到了 Uniffle 场景下虽然不会直接导致节点 OOM但会给某个特定的 Shuffle Server 带来过高的负载。我这里有一套自己的处理思路。首先在应用层通过加盐或者二次分区的方式让数据分布更均匀然后在 Uniffle 层增加分区内多个服务器的冗余存储也就是同一个分区的数据同时散落到多台 Server最后在 Coordinator 层把分配策略从IO调整为IO加Memory加权的方式。三层一起做效果比较明显。另外如果业务上能接受的话把 Spark 的spark.sql.shuffle.partitions调大是一个治本的办法。分区多了单分区的数据量自然下降Server 的负载也就分散了。5.3 性能调优的核心参数组合我在生产环境经过多轮压测总结了一套相对稳定的参数组合可以作为初始配置参考参数推荐值调优说明rss.server.memory.shuffle.highWaterMark.percentage60超过后强制刷盘rss.server.memory.shuffle.lowWaterMark.percentage40刷盘后低于该水位停止强制刷盘rss.server.read.buffer.capacity16m服务端单次读取缓冲区rss.client.write.buffer.size64m客户端内存缓冲调大减少网络请求次数rss.client.read.buffer.size16m客户端读取缓冲区spark.shuffle.io.maxRetries10减少偶发网络故障导致的失败spark.shuffle.io.retryWait5s重试等待时间值得强调的是客户端写缓冲区不能盲目调大。因为它是每个 Executor 的每个 Shuffle Server 连接上都会分配的如果spark.rss.client.type是GRPC那内存开销 连接数 × 缓冲区大小。我最初把写缓冲区调到 128m结果高峰期直接把 Executor 内存打爆了后来调回 64m 才稳定。6. 接入后的稳定性变化与收益分析Uniffle 上线到现在已经有几个月时间这期间的收益数据挺有说服力的。拿我们最大的一个 Spark 生产集群举例单日跑 2000 多个作业高峰期同时运行 300 个作业Shuffle 数据量峰值能达到几十 TB。接入前NodeManager 因为磁盘 IO 过高导致的容器失败每周都会发生好几次几乎每次都需要重跑作业平台稳定性评分一直不达标。接入后NodeManager 上不再有 Shuffle 写入磁盘 IO 降了大约 70%容器失败次数基本归零。虽然 Shuffle Server 也会有故障但因为 Coordinator 会及时摘除故障节点并重新分配对作业的影响被大大降低了。另外还有一个意外收获因为 Shuffle 数据不再保存在计算节点本地节点的磁盘规格需求大幅降低我们成功把一批计算节点的磁盘从 4T SSD 降到了 1T 普通盘硬件成本节省了不少。这对于预算有限但又要跑大作业的团队来说是一个很实际的考量点。7. 到底要不要引入 Uniffle我的个人体会最后分享一点个人经验。很多团队在考虑是否引入 Uniffle 时会纠结于多一套集群是否值得这个问题。从我的实际运行经验来看如果你的集群满足以下任一条件Uniffle 就值得认真评估Shuffle 数据量大且峰值明显作业稳定性要求高容器失败超过 2% 就会引发关注有弹性扩缩容需求Shuffle 占用的本地磁盘空间影响到了计算节点的正常调度。反过来说如果你们的作业都是小数据量、低并发、Shuffle 不频繁那原生的 Shuffle 方式可能已经够用引入 Uniffle 反而会增加运维复杂度。技术选型这件事最忌讳的就是盲目跟风。一个比较务实的验证方法是先在一套测试集群上部署 Uniffle选取几个典型的 Shuffle 密集型作业做对比测试对比作业耗时、失败率、节点 IO 这三个核心指标用数据说话。我当时就是这么做的一周内就坚定了引入的决心。另外有一点需要提前有心理准备从传统的改 NodeManager 参数思路转变为独立 Shuffle 集群思路团队内部需要一定的适应期。尤其是运维同学要多花点时间熟悉 Coordinator 的调度逻辑和 Shuffle Server 的监控指标才能在故障时快速定位问题。我建议从一开始就把 Grafana 监控做起来核心指标包括 Shuffle Server 的读写吞吐、队列长度、GC 耗时和磁盘使用率这些数据在关键时刻能救命。总的来说Uniffle 是我近几年接触过的大数据组件里架构设计思路最清晰、落地收益最直接的项目之一。它用一个相对简单的模型解决了困扰大数据工程师多年的 Shuffle 稳定性问题。如果你正在为 Shuffle 性能发愁花点时间研究一下 Uniffle大概率不会让你失望。
返回列表