Spark YARN集群利用率不足及任务缓慢问题排查咨询
问题1:为何仅使用50个工作节点中的12个?
核心原因是你的任务并行度被RDD分区数直接限制了:
- 你创建的RDD是72个分区,这意味着Spark最多同时运行72个任务(每个分区对应一个独立任务)。
- 你的executor配置是每个executor分配3核,而Spark中1核对应1个并行任务,所以每个executor可以同时处理3个任务。要跑满72个任务,只需要
72 ÷ 3 = 24个executor就足够了。 - 你规划每个工作节点跑2个executor,24个executor刚好需要
24 ÷ 2 = 12个节点。剩下的38个节点因为没有任务可执行,自然处于闲置状态。
简单来说:Spark不会启动超过任务需求的executor,你的分区数直接决定了最大并行任务数,进而限制了实际使用的节点数量。
问题2:为何任务运行如此缓慢?
这个问题涉及资源配置、数据处理、任务特性等多个环节,我拆解几个关键原因:
1. Gensim任务的单线程特性与资源不匹配
你提到串行单任务需要3小时,而gensim的大部分NLP操作(比如词向量训练、文本预处理)都是单线程的。虽然你给每个executor分配了3核,但每个任务只能用到1核——这意味着每个executor的另外2核其实是闲置的,不仅没提升效率,反而可能因为executor之间的内存竞争、线程上下文切换带来额外开销。
2. 数据读取的IO瓶颈
每个分区需要从Blob存储读取两个文件,这里可能存在几个坑:
- 小文件过多:如果文件体积很小,任务的大部分时间会消耗在文件打开、关闭、网络传输上,而非NLP计算本身。比如单文件只有几KB,读取开销远大于处理时间。
- 数据本地性差:如果Blob存储的文件块没有和任务运行的节点在同一区域/可用区,Spark需要跨节点远程读取数据,大幅增加IO耗时。
- 存储吞吐量限制:如果用的是标准层级Azure Blob,大量并行读取可能触发限流,拖慢数据加载速度。
3. YARN内存配置遗漏了堆外开销
你设置了--executor-memory 19G,但YARN模式下executor的总内存包含堆内存和堆外内存(Overhead)。默认堆外内存是executor内存的10%(仅1.9G),但gensim处理NLP任务时,加载预训练模型、处理大文本会用到大量堆外内存,这会导致:
- 频繁的GC(垃圾回收),占用大量CPU时间;
- 甚至触发YARN内存超限,导致executor被杀死重启,进一步拖慢任务。
4. 任务重复加载资源(比如Gensim模型)
如果你的代码在map函数里加载gensim模型,每个分区任务都会重复加载一次模型。比如一个模型加载需要10分钟,72个任务就要消耗720分钟的总加载时间,这会极大拖慢整体进度。
5. 数据倾斜或任务不均衡
虽然你有72个分区,但如果某些分区对应的两个文件远大于其他分区,这些大分区的任务会耗时更久,成为整体任务的瓶颈——其他任务都完成了,只剩几个大任务在跑,导致整体进度缓慢。
关键配置优化建议
针对你的场景,给你几个具体的调整方向:
匹配executor配置与任务特性:
- 把
--executor-cores改为1,同时设置--num-executors 72(刚好匹配分区数),每个节点跑6个executor(留2核给系统),这样每个executor的资源完全给单个任务,避免核心闲置和上下文切换。 - 若想利用所有50个节点,可以把RDD分区数调整为
50节点 × 2executor/节点 ×3核/executor=300个(前提是数据可以拆分到更多分区),这样能充分调动所有节点的资源。
- 把
优化YARN内存配置:
- 添加
--conf spark.yarn.executor.memoryOverhead=4G(建议设为executor内存的20%-30%),确保堆外内存足够,减少GC和内存超限问题。 - 可以适当降低堆内存,比如
--executor-memory 16G --conf spark.yarn.executor.memoryOverhead=8G,总内存24G/executor,两个executor占用48G,节点剩余16G给系统,完全足够。
- 添加
优化数据读取:
- 合并小文件:把多个小文件合并成和HDFS块大小匹配的大文件(比如128M),减少文件读取开销。
- 确保Blob存储的文件块和Spark节点在同一区域/可用区,提升数据本地性。
- 若使用ADLS Gen2,可切换到Premium层级提升吞吐量。
优化Gensim资源加载:
- 使用
SparkContext.broadcast()广播gensim模型,让每个executor只加载一次,而非每个任务重复加载,节省大量时间。
- 使用
排查数据倾斜:
- 统计每个分区的文件大小和处理时间,找出耗时最长的分区,调整分区策略(比如按文件大小拆分大分区)。
内容的提问来源于stack exchange,提问作者xbot

