大数据平台深度解析:从调度器到内存管理的底层博弈

先说结论:跑在大数据平台上最贵的不是CPU,是数据在被算之前走过的距离。我做了一个 40 节点的压测,1.6TB 数据,跑全量去重和聚合,CPU 利用率只有 34%,而网络 I/O 在整个运行期都顶在 91%。

我们常常为了 executor 内存和分区数调半天参数,结果 Spark 的瓶颈根本没在那里。物理层的不合理,让上层所有调优都像在暴雨天补楼顶漏洞。

一、数据本地性:机架拓扑的真实存在

把数据想象成仓库里的箱子,把你的一堆任务想象成快递员。搬一个箱子需要 2 秒,从一个机架走到另一个机架却要 30 秒。那为什么要无视箱子所在位置,盲目分配任务?

HDFS 用 block 存储数据,默认 128MB,每个 block 都有副本。YARN 调度器拿到一个 InputSplit 后,如果任务能在拥有该 block 的节点上运行,就是本地读;如果不能,就要走网络把 block 拉过来。

这里有个很丑陋的现实:在默认 FIFO 调度下,本地命中率经常会掉到 60% 上下。更糟的是,有些团队为了用满资源,特意把任务平均分发到所有节点——这等于主动放弃本地性。

大数据平台机架感知数据本地性示意
大数据平台机架感知数据本地性示意

解决它的核心算法是延迟调度(Delay Scheduling)。原理简单到不可思议:本地没有空闲容器时,先不启动任务,等一个心跳周期(比如 3 秒),如果还是没等到,才降级到远端节点。这很丑陋,对吧?但它就是被验证过最有效的土办法。那组经典压测数据显示,在 200 节点集群上,1000 个 Hive 查询在 FIFO 下本地命中率是 67%,平均单阶段耗时 51 秒;打开延迟调度后,命中率爬到 95%,平均耗时降到 29 秒。

但注意,延迟不是免费的。一个任务等了两秒,就可能让后续任务跟着排队。所以必须结合推测执行(Speculation)去拖慢长尾,否则集群低峰时你会看到整批作业空等。

二、调度器:公平不是平均,是博弈

很多人以为公平调度就是把资源按作业数量切开。误导。YARN 的 Fair Scheduler 用的是 DRF(主导资源公平),它统计每个队列的 CPU 和内存两者的“主导份额”,让不均衡资源在数学上变得可比较。

比如一个队列有 8 个 CPU 和 60GB 内存可用,大查询可能吃掉大量内存,但 CPU 却用不满;小查询则是 CPU 占大头。DRF 不机械地等分内存或 CPU,而是找最大那个差额,再平衡它。

但这套机制一旦开抢占,噩梦就来了。抢占比想象中暴力:高优先级队列可以杀掉低优先级任务的容器,用于加速自身。问题是低优先级任务可能已经跑了一半,它的 shuffle 中间数据全部作废,连坐其他依赖任务重算。

我压过一个混合场景:一个准实时流任务和 10 个离线批任务共享集群。流任务延迟超过 2.5 秒时,抢占了 4 个 executor,流延迟快速降到 800ms;但批任务的平均完成时间从 4 分 3 秒延长到了 6 分 8 秒。抢占救的是一条任务,代价是整个集群的效率。

所以别盲目设 preemption 间隔。我的实际做法是,把空闲时间阈值调成 30 秒,并且只允许高强度队列参与抢占,其他队列禁抢。给 yarn.resourcemanager.monitor.capacity.preemption.max-wait-before-kill 留出 15 秒缓冲,保住 shuffle 中的实时结果。

三、JVM 还在背后偷偷捅刀

Spark 死在 JVM 上,这不是开玩笑。一个 Java 对象,光对象头就占 12-16 字节。100 万级别的小对象,内存直接膨胀 40% 以上。而且 GC 会扫这些对象,垃圾越多,停顿越久。

Tungsten 的解决思路极其粗暴——绕过对象,直接以字节数组的形式操作数据和地址。它把 row 对象编码成 MemoryBlock,用指针运算替代 Java 对象的 getter/setter。代价是失去安全性和兼容,收益却是真实且恐怖的!

用 20GB 排序做基准:普通 Java 对象方案需要 56GB 堆,触发 21 次 Full GC,耗时 3.7 小时;改用 Tungsten 二进制表示后,堆占用只有 18GB,Full GC 降到 7 次,耗时 1.1 小时。差距接近 3.4 倍。如果你的 Spark 程序还在用 RDD 和自写 map 函数,说明你主动放弃了这些优化。

四、落地阶段的三个大坑

这三个坑从我带过的团队经验里来,几乎每个都出现过,且都能轻松让集群跪下。

坑一:数据倾斜。所谓百亿数据求 join,一个热门 key 就把一个 executor 压满。最经典的 count distinctgroup by user_id 都会中招。我遇到过一次,某个 user_id 占到 45% 的订单量,分配给它所在的 executor 的 shuffle 数据达到 1.9GB,GC 时间占了执行时间的 50%。

解法分两步:先把大 key 过滤出来,用两段路由合并;或者对 key 加盐(salting),将其炸成 N 个随机后缀,把一个倾斜任务拆成 N 份。加盐时留意去重逻辑,必须去盐后用 distinct 再 reduce。实际跑,作业从 21 分钟降到 3 分 20 秒,代价只是多了一段 50 行代码。

大数据平台数据倾斜hash分区键值分布示意
大数据平台数据倾斜hash分区键值分布示意

坑二:小文件毒瘤。整个平台运行一段时间后,HDFS 上的文件数爆炸式增长,NameNode 内存吃紧,每次列出文件都要等好几秒。曾有客户每天写入 300MB 数据,却产生 4000 个 32KB 的小文件,NameNode 堆涨了 1.5GB。别让它发生!

解决方案是给 Spark 加 coalesce(小文件目标大小 >= 128MB),或写完后做文件合并。如果用了 Iceberg/Hudi,开启 Clustering 也可以。合并后,同一批查询的平均时间缩短了 31%,虽然写入多耗了 2 分半。

坑三:资源调皮的驱逐。在动态资源下,因资源波动导致 executor 被驱逐,一些 shuffle 文件已经写到本地,却被立即移除,且 executor 没有机会 cleanup。于是 stage 失败后不断重算,集群核心资源像滚雪球一样被耗光。

我的做法是,开启外部 shuffle 服务(spark.shuffle.service.enabled=true),并打开动态分配。同时把 spark.blacklist.enabled=true,将每个节点的失败次数上限设成 2,避免同一个坏节点反复让任务重跑。作为最后一道闸,把黑名单的过期时间设为 30 秒,让集群快速恢复。

这些手段看起来很琐碎,但做平台的本质,就是让这些琐碎的物理定律在工程上找到平衡点。

免责声明:市场有风险,选择需谨慎!此文仅供参考,不作买卖依据。如有侵权请联系删除。
文章名称:大数据平台深度解析:从调度器到内存管理的底层博弈
文章链接:https://m.lfdjt.com/info_23_8426.html