如何确认Spark应用是否充分利用集群所有可用资源?
如何确认Spark应用是否充分利用集群所有可用资源?
兄弟,你的情况我太有共鸣了!一开始看到Executor日志全是GC信息,确实会慌,以为集群资源都浪费了,但其实不一定——咱们一步步拆解排查,搞清楚到底是真没干活,还是只是日志没显示而已。
首先得澄清一个误区:Executor默认不会输出你的业务代码日志!除非你在代码里专门配置了日志框架(比如Log4j2)把业务日志输出到控制台,否则Executor的日志里只会有Spark本身的系统日志(比如GC、Task调度信息)。所以只看到GC日志≠Executor在摸鱼,得用更靠谱的方式验证。
下面是几个关键的排查步骤,都是我平时调Spark应用常用的:
一、直接看Spark UI(最权威的方式)
你在YARN Resource Manager里找到你的应用,点进去ApplicationMaster的链接就能打开Spark UI,重点看这几个页面:
- Jobs页面:看每个Job的Stage数量和执行状态,尤其是读取Kafka的那个初始Stage——它的Task数量应该等于你所有Kafka主题的分区总数(你50个主题每个至少4分区,那至少200个Task)。如果Task数远低于这个数,那说明并行度没拉起来。
- Executors页面:看每个Executor的「Active Tasks」「Completed Tasks」数值,还有CPU、内存的使用率。只要Completed Tasks大于0,就说明这个Executor干过活;如果某个Executor的Active Tasks一直是0且Completed Tasks为0,那才是真的没分配到任务。
- Stages页面:点进读取Kafka的那个Stage,看Task的分布情况——是不是所有Executor都有Task在跑,有没有某个Executor从头到尾没接过任务。如果Task集中在少数几个Executor上,那就是资源没充分利用。
二、检查Kafka数据源的并行度匹配
Spark读取Kafka时,默认的Task数等于Kafka主题的总分区数,这个是保证并行读取的关键。你要确认:
- 代码里有没有对Kafka读进来的DataFrame做
repartition()或者coalesce()把分区数改少了?如果有,那会直接降低并行度,导致很多Executor闲下来。 - 你的Kafka主题分区是不是真的每个都有4个?可以用Kafka命令行工具确认:
kafka-topics.sh --describe --topic <你的主题名> --bootstrap-server <Kafka地址>。
三、检查你的Spark提交配置
你贴的spark-submit命令没写完,我猜可能没配置Executor的核心参数,这在EMR上很容易导致默认资源不够:
- 有没有设置
spark.executor.instances(Executor数量)、spark.executor.cores(每个Executor的CPU核数)?比如如果默认每个Executor只有1核,那即使有10个Executor,也只能并行跑10个Task,剩下的Task都得排队,看起来就像很多Executor没干活。 - 建议你补上这些配置,比如根据你的集群规模调整:
spark-submit --name read-from-kafka --deploy-mode cluster --master yarn \ --conf spark.eventLog.enabled=false --conf spark.sql.caseSensitive=true \ --conf spark.sql.shuffle.partitions=50 --conf spark.driver.memory=5300M \ --conf spark.executor.instances=10 --conf spark.executor.cores=4 \ --conf spark.executor.memory=8G \ --class com.XXX.XXX.reports.unifiedLoader --jars file:////home/hadoop/lib/* ...
四、排查任务瓶颈
如果Spark UI显示Task确实分配到了所有Executor,但还是觉得效率低,那可能是任务本身有瓶颈:
- 比如写入数据库的阶段,如果DB的写入吞吐量有限,那后面的Task会排队,导致部分Executor暂时空闲,但这不是集群没利用,而是下游系统拖了后腿。
- Executor里GC日志多,也可能是内存配置不合理——比如
spark.executor.memory设太小,导致频繁GC,即使在干活效率也低,这时候可以调大内存或者调整内存比例(比如spark.executor.memoryOverhead)。
最后给你个小技巧:如果想在Executor日志里看到业务日志,你需要在代码里配置Log4j2,把日志级别设为INFO,并且输出到控制台,这样YARN就能把业务日志捕获到Executor的日志文件里了。
备注:内容来源于stack exchange,提问作者mt_leo
相关产品推荐
相关产品推荐

