Dataproc Spark集群YARN CPU使用率与htop数据不一致问题排查
问题:Dataproc YARN监控CPU使用率与htop实际占用不符的原因及优化方法
集群环境与配置
- 操作系统:Ubuntu 18.04
- Spark版本:3.3.0
- 集群节点配置:
- Master节点:内存7.5GiB,CPU核心数2,主磁盘32GB
- Worker节点:CPU核心数16,内存16GiB,YARN可用内存13536MiB,主磁盘32GB
SparkSession初始化代码
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.functions import udf from pyspark.sql.types import StringType
spark = SparkSession.builder.\ config("spark.executor.cores","15").\ config("spark.executor.instances","2").\ config("spark.executor.memory","12100m").\ config("spark.dynamicAllocation.enabled", False).\ config("spark.sql.adaptive.enabled", False).\ config("spark.sql.files.maxPartitionBytes","10g").\ getOrCreate()
数据与任务情况
- 读取并缓存约40GiB的CSV文件,生成30个分区;缓存后内存中反序列化数据10.3GiB,磁盘序列化数据3.9GiB。
- CPU验证任务代码:
@udf(returnType=StringType()) def f1(x): out = '' for i in x: out += chr(ord(i)+1) return out @udf(returnType=StringType()) def f2(x): out = '' for i in x: out += chr(ord(i)-1) return out df_covid = df_covid.withColumn("_catted", F.concat_ws('',*df_covid.columns)) for i in range(10): df_covid = df_covid.withColumn("_catted", f1(F.col("_catted"))) df_covid = df_covid.withColumn("_catted", f2(F.col("_catted"))) df_covid = df_covid.withColumn("esize1", F.length(F.split("_catted", "e").getItem(1))) df_covid = df_covid.withColumn("asize1", F.length(F.split("_catted", "a").getItem(1))) df_covid = df_covid.withColumn("isize1", F.length(F.split("_catted", "i").getItem(1))) df_covid = df_covid.withColumn("nsize1", F.length(F.split("_catted", "n").getItem(1))) df_covid = df_covid.filter((df_covid.esize1 > 5) & (df_covid.asize1 > 5) & (df_covid.isize1 > 5) & (df_covid.nsize1 > 5))
- 现象:执行
df_covid.count()触发计算后,htop监控显示两个Worker节点的16核均被完全占用(持续3-4分钟),但Dataproc的YARN监控显示CPU使用率最高仅约70%。
原因分析
统计口径差异
- YARN仅统计自身管理的容器(Spark Executor的JVM进程)的CPU使用率,不包含节点上系统进程、Dataproc后台服务(如监控Agent、日志采集进程等)的CPU占用;而htop统计的是节点所有进程的总CPU使用率,当Executor满负荷时,加上后台服务的占用会让htop显示100%,但YARN只统计容器部分,数值自然更低。
- 此外,YARN默认基于虚拟核计算使用率,若Worker节点开启超线程,YARN统计的是虚拟核的使用情况,而htop显示的是物理核的实际占用,二者维度不同也会造成数值偏差。
Executor CPU配置未占满节点核
- 当前配置给每个Executor分配15核,每个Worker节点运行1个Executor(共2个Executor对应2个Worker),但Worker节点有16核,剩余1核被YARN NodeManager等系统进程占用。YARN统计时仅计算Executor的15核使用情况,单节点YARN理论最高使用率为15/16≈93.75%,再加上任务调度间隙、GC停顿等实际损耗,最终显示70%左右的使用率符合预期。
Python UDF的跨进程开销
- Python UDF需要在JVM与Python子进程间进行数据序列化/反序列化,这部分额外开销会稀释JVM端的CPU使用率。YARN监控的是JVM进程的CPU,而htop包含了Python子进程的CPU占用,这进一步拉大了二者的数值差异。
优化方案
1. 调整Executor CPU配置,占满节点物理核
将spark.executor.cores调整为16,保持spark.executor.instances为2(每个Worker运行1个Executor)。让YARN容器占用Worker节点全部物理核,消除预留核的损耗,提升YARN统计的CPU使用率上限。
2. 替换Python UDF为Spark内置函数
Python UDF的跨进程开销是性能瓶颈之一,将自定义的f1、f2替换为Spark内置的字符串操作函数,避免JVM与Python进程的交互开销,提升JVM端CPU使用率,让YARN监控更接近实际节点CPU占用:
df_covid = df_covid.withColumn("_catted", F.concat_ws('',*df_covid.columns)) # 利用Spark内置translate函数实现字符偏移逻辑 from pyspark.sql.functions import translate # 生成ASCII字符映射表 lower_chars = "abcdefghijklmnopqrstuvwxyz" upper_chars = lower_chars.upper() full_chars = lower_chars + upper_chars # 生成偏移后的映射表(+1和-1可以通过反向映射实现) shifted_plus1 = lower_chars[1:] + lower_chars[0] + upper_chars[1:] + upper_chars[0] for i in range(10): df_covid = df_covid.withColumn("_catted", translate(F.col("_catted"), full_chars, shifted_plus1)) df_covid = df_covid.withColumn("_catted", translate(F.col("_catted"), shifted_plus1, full_chars))
内置函数在JVM内执行,没有跨进程开销,能让CPU使用率更集中在YARN监控的容器进程中。
3. 调整分区数,匹配Executor并行度
当前数据有30个分区,总Executor核数为15*2=30,刚好匹配,但可以尝试将分区数调整为核数的1.5-2倍(即45-60个分区),避免个别分区过大导致的负载不均,提升整体CPU利用率。可以通过df_covid.repartition(45)或者调整spark.sql.files.maxPartitionBytes到5GiB左右来实现。
4. 精简节点后台进程
检查Worker节点上的Dataproc后台服务,关闭非必要的日志收集、监控Agent等进程,减少系统进程占用的CPU,让YARN容器占用更高比例的节点CPU,从而提升YARN监控显示的使用率。
内容的提问来源于stack exchange,提问作者figs_and_nuts
相关产品推荐
相关产品推荐

