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

资讯详情

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

Apache Uniffle实战:彻底解决Spark Shuffle性能瓶颈与稳定性问题

Apache Uniffle实战:彻底解决Spark Shuffle性能瓶颈与稳定性问题 在排查当天凌晨几个跑批的任务时Spark SQL 的 shuffle 阶段又成了背锅侠——磁盘吞吐被打满executor 老是被判定为超时整个计算队列一片哀嚎。这类问题很多做过大数据平台的人都懂shuffle 明明是计算引擎里最绕不开的环节也是最容易把集群拖垮的环节。后来我把这些任务的 shuffle 链路整个迁到了 Apache Uniffle 上效果立竿见影任务稳定性和整体的 IO 压力都改善了很多。Apache Uniffle 是个开源的统一 Shuffle 引擎主要解决 Spark、MapReduce 等计算框架在 shuffle 阶段碰到的一系列性能瓶颈和可靠性问题。它从腾讯内部的大数据实践里剥离出来后来捐献给 Apache 基金会完成了孵化现在已经是社区里做 remote shuffle service远程 Shuffle 服务比较有代表性的组件之一。这篇文章我不会通篇堆概念而是按照我自己评估、部署、接入、排障的完整过程把 Uniffle 的原理、架构、安装配置、踩坑经验一次讲透给正在被 shuffle 问题折磨的同学一份能直接抄作业的参考。顺便回应一下搜索框里的热词knuth shuffle 里的科努特的确是位数学家全名 Donald E. Knuth也就是计算机领域无人不知的高德纳。他写的那套《计算机程序设计艺术》对整个算法领域影响极深Fisher-Yates 洗牌算法被他在书中以严谨的数学方式整理描述之后很多人就把这种洗牌算法叫做 Knuth Shuffle。这里的 shuffle 与我们说的大数据 shuffle 在核心思想上有一点相通——都是要把数据按照某种规则打散、重排只是前者面向数组随机排列后者面向分布式计算中的海量数据重新分区。理解了这一层再来看 Uniffle 要解决的问题思路会清晰很多。1. Shuffle 的本质与原生方案的痛点1.1 为什么每个分布式计算引擎都绕不开 Shuffle在分布式计算里数据大概率分散在多台机器的多个执行进程里。比如一个 Spark 任务处理 1 TB 的日志输入数据被切分到几百个分区里分散在不同机器上并行计算。如果下一个计算阶段需要按照某个 key比如用户 ID把相同 key 的数据汇聚到同一个节点继续处理那么数据就必须跨节点重新分配——这个重新分区、跨节点搬移数据的过程就是 shuffle。这个过程很像“发快递”。map 端上游处理完的节点是各个发货点reduce 端下游需要聚合数据的节点是各个收货点。发货点需要把自己手头的货按照目的地的不同分拣、打包然后一件一件运出去收货点要把从不同发货点来的货统一签收、整理再开始下一轮工作。任何一个环节没规划好比如某个发货点分拣太慢、某个收货点处理的包裹远多于别的收货点整个物流体系就卡住了。放到实际计算框架里shuffle 就是 map 阶段到 reduce 阶段之间的那一段数据流转。它对任务的完成时间、集群资源的利用率、甚至任务的成败都有决定性的影响。很多任务看起来很复杂但主要的执行时间大头恰恰不是计算逻辑本身而是 shuffle 阶段的数据写入、传输、读取。1.2 原生 Shuffle 有哪些让人头疼的问题先说说原生 shuffle 是怎么做的这样才更能体会到 Uniffle 改了哪些东西。以 Spark 的 Hash Shuffle 和 Sort Shuffle 为例map 端把处理完的数据按 key 分区后通常先写到自己所在节点的本地磁盘上。下游的 reduce 任务开始执行时再直接到对应的上游节点上拉取属于自己的那部分数据。这套方案看着简单实际运行起来问题不少。第一上游节点本地磁盘压力巨大。shuffle 数据要落盘而且写的是临时文件用完就删。磁盘带宽和 IO 吞吐成为瓶颈。我曾经遇到过一台机器上同时跑了多个 executor每个 executor 都在疯狂写 shuffle 临时文件结果整块数据盘的 IO 被打到 95% 以上连系统监控都开始延迟报警。这种情况在小集群里尤其致命因为一个节点的磁盘抖动会影响该节点上运行的所有任务。第二节点故障导致数据丢失。如果上游节点在 shuffle 数据还没被下游全部拉走时就挂掉了那部分数据就彻底丢失。Spark 的解决之道是把上游任务重新计算一遍重新生成 shuffle 数据也就是 stage 重试。问题是大数据任务里一个 stage 的计算往往很重跑了一个多小时的数据因为一个节点宕机就要全部重来成本高得让人肉疼。第三数据倾斜加剧单点瓶颈。某个 key 的数据量和别的 key 不在一个量级时对应的下游节点会拉取远超预期的数据量。在原生 shuffle 架构下这些数据全部拥挤在一台机器的有限磁盘、网络和 CPU 资源上任务跑得慢只是开始更惨的是直接 OOM 或者被判定为失效。第四资源调度毛刺。各个 mapper 完成时间参差不齐reduce 端为了拉数据需要反复试探等待如果某些 mapper 卡住整个 stage 的完成时间就被无限拉长。这种现象在线上我见过太多明明几十秒就能跑完的 stage硬生生被拖到十几分钟。1.3 从“进程内洗牌”到“独立服务化”原生 shuffle 是在计算进程内部完成的每个 executor 既是生产者也是存储者既是消费者也是传输者。这种高度耦合的架构在任务规模小的时候没什么问题一旦集群规模大了、任务复杂了所有问题都被放大。于是业界开始思考能不能把 shuffle 这个环节从计算框架里抽出来做成一个独立服务这就是远程 Shuffle 服务的思路。Uniffle 就是这个思路的践行者之一。它把 map 端生成的数据先发送到一个独立的 Shuffle Server集群由这些专门的服务器来接收、合并、存储、管理 shuffle 数据。reduce 端要数据时不再直接去一个个上游节点拉取而是统一从 Shuffle Server 集群获取。这样一来计算节点的本地磁盘压力大大缓解节点故障导致的数据丢失问题也因为多副本机制得到改善整个 shuffle 链路走向了服务化和专业化。本质上Uniffle 做的是“把分布式计算里的洗牌行为从进程内部搬到独立服务”让做计算的专心算让管数据的专心管。下面我详细拆解一下它的架构和核心设计。2. Uniffle 整体架构与核心设计2.1 两大服务端角色Coordinator 与 Shuffle ServerUniffle 的架构很清晰从服务端角色的视角看主要就是两个角色一个管调度一个管数据。Coordinator 相当于调度中心。它负责维护集群里所有 Shuffle Server 的存活状态、资源情况和数据分布元信息。计算引擎需要为一个 shuffle 阶段分配资源时就会向 Coordinator 申请一组可用的 Shuffle ServerCoordinator 根据当前集群负载和可用资源情况返回一个合适的服务节点集合。同时Coordinator 还负责容错处理比如某个 Shuffle Server 宕机了Coordinator 需要感知到并在后续分配时尽量绕开故障节点同时也要处理因为节点故障导致的数据恢复问题。Shuffle Server 是数据面节点。它接收 map 端写过来的 shuffle 数据在内存和本地磁盘或 HDFS上组织存储等到 reduce 端来读取时再把数据吐出去。每个 Shuffle Server 会在启动时注册到 Coordinator 上并周期性上报心跳包含当前磁盘剩余空间、内存使用情况、正在处理的任务数等负载指标让 Coordinator 在分配时可以参考。这里要注意Uniffle 里一个完整的 shuffle 任务叫一个 shuffle提交过程中会涉及到 application、shuffle id、partition 等几个层级的概念。简单理解一个 Spark stage 对应一批 shuffle 数据每个 shuffle 数据按分区组织Shuffle Server 在存储时就是按 “应用-任务-分区” 的维度来管理这些数据块。2.2 数据写入链路从 Mapper 到 Shuffle Server当我用 Spark 接入 Uniffle 之后一个 stage 的执行过程会变成这样。map 端每个 task 处理完自己的数据之后不再往本地磁盘写 shuffle 文件而是先把数据在内存中按目标分区进行缓存和合并然后通过 RPC 发送给 Coordinator 分配好的 Shuffle Server 节点。这些数据到达 Shuffle Server 后服务端会先把数据缓冲到内存里达到一定阈值后批量地刷入底层存储。因为 Shuffle Server 专门干这件事它的数据刷写策略可以做得更加激进和高效比如说做更多的内存合并、更大的批量写让落盘的数据块更整洁、更连续。这个批量写的优化很重要。原生 shuffle 每个 map 会为每个 reduce 生成一个文件片段结果就是大量零碎的小文件。而 Uniffle 在服务端把属于同一个分区、来自不同 mapper 的数据合并成大的数据块再落盘文件数量少、大小整齐读取时磁盘寻道开销也小很多。写入链路里还有一层非常关键的机制就是数据合并。map 端输出的数据是按行的如果每条数据都单独发一次 RPC网络包数量会爆炸。Uniffle 会把属于同一个分区的多条数据攒成一个更大的 buffer达到一定大小之后再发送。这个攒批的过程会引入一点延迟但对整体吞吐的提升是巨大的。实际配置中有个参数控制单个缓冲区的大小单条数据特别大的任务和大量小消息的任务需要差异化调优。2.3 数据读取链路Reducer 如何获取数据reduce 端读取数据的方式也变了。在原生 shuffle 中reducer 需要同时对上几百个甚至上千个 mapper 发起连接请求试想一下几千个 reducer 同时对几千个 mapper 发起连接整个集群的网络连接数是百万级别光连接的管理和等待就能拖垮节点。Uniffle 下reducer 只需要连接集群里分配好的几个 Shuffle Server 节点而且它会批量地从一个节点上一次性读取整个分区数据而不是按 mapper 挨个去要。以 Spark 为例reduce 端在读取时会先从 Coordinator 查询自己需要的分区数据在哪些 Shuffle Server 上然后并行地从这些服务端拉取数据。服务端读取时也是顺序读减少了随机 IO 的次数。还有一个对外行感知比较明显的差异——reduce 端在原生 shuffle 里会因为等待某个 mapper 的数据而卡住整个 stage 都不能完成。Uniffle 因为数据已经全部上传到了服务端mapper 执行完就完事了reduce 端拉取数据不依赖任何 mapper 的存活状态。所以只要数据完整性有保障整个 stage 的推进速度不会因为个别节点掉队而陷入停滞。2.4 可靠性设计为什么要做多副本和故障恢复Uniffle 在可靠性上投入很大这也是它区别于很多“半成品”远程 Shuffle 实现的重要一点。数据从多个 mapper 上传到 Shuffle Server 后如果这个 Server 突然宕机那这些数据就危险了。Uniffle 提供多副本机制也就是每份 shuffle 数据同时发送给多个 Shuffle Server 保存默认情况下可以配置为 2 副本或者 3 副本。只要还有任意一个副本存活数据就不会丢下游任务就可以正常读取完全不需要上游重算。这个机制对长任务的可靠性提升非常明显。我接入之前最担心的就是这个毕竟把数据从计算节点挪到了独立的服务节点如果服务节点的稳定性还不如原来的 executor那不反而更糟实际上 Uniffle 的多副本机制、Coordinator 的调度容错机制都已经比较成熟节点挂掉后Coordinator 能在很短的时间内把故障节点上正在服务的数据重新分配、恢复校验让任务继续进行下去。当然代价也很直接多副本意味着双倍的存储空间和写入流量。所以生产环境到底配 2 副本还是 3 副本需要根据数据的业务重要性、集群规模、磁盘成本综合判断。一般来说日志分析这类可以容忍一定程度重跑的任务2 副本够用了重量级的离线报表、核心数据加工链路强烈建议 3 副本。3. 部署实践与核心配置3.1 规划集群部署形态与硬件侧重Uniffle 的部署形态比较明确一个集群包含若干 Coordinator 节点和若干 Shuffle Server 节点。Coordinator 通常部署 2 到 3 个做高可用Shuffle Server 则根据数据量和并发度扩展。Coordinator 节点对硬件要求不高主要是管理元数据和做调度决策8 核 16 GB 内存就非常充裕了。Shuffle Server 是数据面节点它对磁盘的要求最高建议每台节点配置多块大容量数据盘内存也要给足因为写入链路里会做内存缓冲合并。我个人的经验是Shuffle Server 的内存可以按总数据盘容量的一定比例来配比如每 10 TB 存储配 32 GB 到 64 GB 内存具体要看单个任务的并发度。Uniffle 支持本地文件系统和 HDFS 两种底层存储。使用本地文件系统时部署最简单性能也最好但数据可靠性会受机器故障影响。使用 HDFS 时天然拥有多副本机制数据可靠性高但多一层网络和 NameNode 的 IO 开销。我个人的建议是先用本地文件系统做验证和 POC线上根据可靠性的要求再决定是否切换到 HDFS。有些团队也用了本地盘搭配 Uniffle 自身的多副本机制这样既保证了性能又保障了可靠性只不过运维成本也会高一点。3.2 配置与启动关键步骤回顾Uniffle 的部署并不复杂。假设我已经准备好了三台机器其中两台做 Coordinator四台做 Shuffle Server每台 Shuffle Server 挂了两块数据盘。先要准备 JDK 环境推荐 JDK 8 或者 11官方对高版本 JDK 的支持也一直在完善但线上稳妥起见我建议用 8 或 11。然后下载对应版本的 Uniffle 二进制发行包解压到统一目录下。Coordinator 的配置在 conf/coordinator.conf 里。有几个关键项rss.coordinator.portCoordinator 的通信端口默认 19999rss.coordinator.app.expired应用元数据过期时间默认可以按需配置rss.coordinator.server.assignment.strategy服务端分配策略常见的有根据负载分配Shuffle Server 的配置在 conf/server.conf 里关键项包括rss.rpc.server.port服务端 RPC 端口默认 19998rss.storage.type存储类型LOCALFILE 表示本地文件系统HDFS 表示使用 HDFSrss.server.buffer.capacity服务端内存缓冲总容量rss.server.read.buffer.capacity读取缓冲区总容量rss.server.flush.thread.alive刷盘的线程数一般跟磁盘数量有关rss.server.disk.capacity每块数据盘可以使用的空间上限rss.server.replica.write写入副本数默认是 2启动方式比较简单Coordinator 使用bin/start-coordinator.shShuffle Server 使用bin/start-shuffle-server.sh。启动完记得用jps或者看日志确认进程真的起来了。Coordinator 起来以后可以访问它的 HTTP 端口查看集群状态Shuffle Server 注册成功与否在上面都能看到。3.3 接入 Spark改动最小化的集成方案接入 Spark 是 Uniffle 使用中最常遇到场景。我需要先确认自己的 Spark 版本和 Scala 版本下载对应的客户端 jar 包比如 uber 包。然后把这个 jar 放到 Spark 的 classpath 下再在 spark-defaults.conf 里加上关键配置。最核心的一段配置长这样spark.shuffle.managerorg.apache.spark.shuffle.RssShuffleManager spark.rss.coordinator.quorumcoordinator-host1:19999,coordinator-host2:19999 spark.rss.storage.typeLOCALFILE spark.rss.client.read.buffer.size16m spark.rss.client.send.buffer.size4m spark.rss.writer.buffer.size4m spark.rss.writer.buffer.spill.size64m这里最关键的一行是spark.shuffle.manager它决定了 Spark 使用哪套 shuffle 管理器。换成 RssShuffleManager 之后Spark 的 shuffle 行为就整体切换到了 Uniffle 上。这里需要特别提醒一下spark.rss.coordinator.quorum里填的一定要是能连通 Coordinator 的地址最好在 Spark 客户端机器上先用 telnet 验证一下端口连通性。另一个容易踩坑的是 executor 的内存里必须给 shuffle 客户端预留一些空间如果 executor 原本内存就压得很紧加了 Uniffle 客户端之后反而容易出现因为缓冲区无法分配而导致的 OOM。全部配置好之后提交一个不算太大的测试任务观察两个指标一是任务日志里是否看到 RssShuffleManager 的初始化信息二是 Coordinator 的 Web 界面或者日志里能否看到任务注册上来的记录。如果这些都正常就说明接入成功后续就可以把线上任务逐步迁移过来。3.4 与原生 Shuffle 的关键参数对比对比原生 shuffle 和 Uniffle 的配置思路你会发现想要的优化点不太一样。原生 shuffle 调参时我要关注的核心参数是spark.shuffle.file.buffer、spark.shuffle.spill.compress、spark.shuffle.memoryFraction这些。调优目标是在有限的 executor 内存里挤出更多空间给 shuffle 缓冲尽量减少溢写到磁盘的次数。而接入 Uniffle 之后executor 本地基本不再承担大的 shuffle 数据存储任务调优重心变成了怎么更高效地把数据推送到 Shuffle Server。需要关注spark.rss.writer.buffer.size这种参数它决定 map 端攒多少数据发一次网络包还要关注服务端的刷盘策略让服务端能够平滑地承接多个 executor 同时写入的压力。这个变化本质上是“把问题从计算节点转移到了服务节点”所以 Uniffle 模式下调优的另一个大头是 Shuffle Server 的资源规划。如果服务节点的磁盘或网络先成为瓶颈再调客户端参数也意义不大。4. 常见问题与排查技巧实录4.1 接入后任务变慢可能败在了哪里这是很多初用者反馈最多的问题明明 Uniffle 宣称能提升 shuffle 性能为什么某些场景下反而变慢了我排查过一个很典型的案例。任务的每条记录特别大而任务并行度又高导致 map 端每条记录在发送缓冲区里只能放很少的条数就达到阈值触发发送。这就出现了一种“频繁小包发送”的模式网络包数量巨大每个包的有效数据占比很低吞吐反而不如原生 shuffle 的批量写盘和批量拉取。解决方案是结合记录大小适当调大发送缓冲区并且合理设置数据合并的逻辑让单次网络包能携带更多有效数据。还有一类情况是任务本身 shuffle 数据量很小比如几十 MB但同一个 SparkContext 里频繁启停 stage。这种场景下Uniffle 带来的额外链路开销反而让总耗时增加了。遇到这种情况我会建议保留部分数据量特别小的短任务走原生 shuffle数据量大的任务走 Uniffle两种模式在同一个集群里混合使用完全没问题。4.2 服务端磁盘或内存成为新瓶颈时怎么办引入 Uniffle 之后瓶颈从“计算节点的本地磁盘”转移到了“Shuffle Server 的磁盘、内存、网络”。这是架构升级带来的正常现象但如果不主动管理和规划服务端很容易被巨大压力击穿。在实际运维中我发现最容易出现问题的是 Shuffle Server 的内存缓冲设置过小导致刷盘频率非常高磁盘 IO 持续高位或者缓冲设置得太大单个大任务写入量一冲上来内存立刻被打满触发频繁的 GC。我的经验是先用监控工具观察 Shuffle Server 的内存使用曲线找出稳定运行时的峰值再按峰值留出 30% 到 50% 的余量来进行设置。磁盘方面Shuffle Server 上保存的临时数据要及时清理。Uniffle 本身有清理机制但清理的触发时机和扫描频率也需要关注。如果任务量特别大磁盘上的过期数据清理不及时可能出现磁盘空间不足。建议配置好磁盘容量上限定期观察磁盘使用趋势避免因为日志和数据清理不及时导致的服务不可用。4.3 常见问题速查清单我把运维过程中碰到过的高频问题整理成一张表方便大家遇到问题的时候快速定位。问题现象可能原因处理方式任务启动即报 Coordinator 连接失败Coordinator 地址错误或端口被防火墙拦截检查rss.coordinator.quorum用 telnet 验证端口连通性任务运行中出现大量 shuffle 读取超时Shuffle Server 负载过高或网络抖动检查服务端 CPU、内存、磁盘 IO确认是否需要扩容或优化分配策略executor 内存暴涨甚至 OOM客户端缓冲区配置过大或 executor 可用内存预留不足减小rss.client.send.buffer.size为 executor 预留更多内存服务端磁盘空间快速打满过期数据清理不及时或副本数配置过高检查清理线程日志合理配置清理周期和副本数读取数据大小和写入数据大小对不上数据在传输过程中出现异常或客户端版本和服务端不一致核对客户端与服务端版本建议升级到同一稳定版本某个 Shuffle Server 宕机后任务仍长时间运行多副本未生效或 Coordinator 容错策略配置不完善检查服务端是否配置了多副本确认 Coordinator 健康检查和故障转移策略生效shuffle 数据量不大但耗时很长并发不足或网络包堆积通过监控工具观察线程池使用情况确认客户端并发连接数配置合理这张表本质上是给自己留底的排查手册每一条背后都是一次真实的线上教训。建议团队在接入 Uniffle 的初期就建立类似的文档每次排障后补充更新后续遇到相似问题能节省很多时间。4.4 选型对比Uniffle 和同类方案怎么选社区里和 Uniffle 定位接近的还有 Apache Celeborn前身是腾讯的 Remote Shuffle Service。两个项目解决的问题类似但各有侧重。Uniffle 在多引擎支持上做得比较全面Spark、MapReduce、Tez 都能接入在海外社区和国内社区都有不错的活跃度。Celeborn 在接口设计和实现上也比较优雅对 Spark 场景的支持成熟度很高也有一些特有的多副本存储策略。从我个人的使用体验来说如果你的集群以 Spark 为主MapReduce 任务也不少希望用一个 Apache 基金会主导、社区持续活跃的项目Uniffle 会是一个比较稳妥的选择。如果你只是解决 Spark 的 shuffle 稳定性和性能问题Celeborn 也可以列入对比测试范围。归根结底选型不能单看文档和特性列表一定要在自己的数据集、自己的任务规模、自己的硬件条件下做一轮小流量压测看真实运行数据和稳定性再做最终决策。5. 一些我坚持使用的调优与运维习惯Uniffle 本身是有生命周期的不断在演进参数也在变化。但有些运维习惯是可以跨版本坚持的。无论任务大小接入 Uniffle 前我会先在测试环境完整跑一遍核心链路并刻意制造一个节点宕机的场景测试多副本和故障恢复是否正常。这个测试过程能让整个团队对 Uniffle 的可靠性建立真实的信心的不至于上线后遇到一点风吹草动就想回滚。再一个是监控。Uniffle 的 Coordinator 暴露了很多指标包括集群分配情况、服务端负载、数据读写吞吐、任务失败率等。接上 Prometheus 和 Grafana把关键指标做成一个大盘比什么都管用。我在生产环境就遇到过磁盘容量分布不均衡的问题如果不看监控大盘光靠任务表现来判断可能要过好久才能发现。最后建议关注社区的版本更新。Uniffle 的版本迭代速度很快不少 bug 修复和性能优化都是通过新版本发布的。我每次升级前都会认真看 release note挑出对自己环境有影响的变更点用测试环境验证后再上生产。Shuffle 这类底层组件最忌讳的是一套配置用三年完全不跟进社区节奏。在折腾 Uniffle 的过程中我最大的感受是它把 shuffle 这个原本隐藏在计算引擎内部的“灰色地带”拉到了聚光灯下让我开始认真审视每个任务的 shuffle 模式、数据量、网络开销和存储行为。一旦这些数据可视化、服务化之后调优才真正有了抓手。如果你也被线上任务的 shuffle 问题困扰不妨照着上面的思路小范围验证一下看看这个统一 Shuffle 引擎能不能给你带来惊喜。
返回列表