You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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()

数据与任务情况

  1. 读取并缓存约40GiB的CSV文件,生成30个分区;缓存后内存中反序列化数据10.3GiB,磁盘序列化数据3.9GiB。
  2. 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))
  1. 现象:执行df_covid.count()触发计算后,htop监控显示两个Worker节点的16核均被完全占用(持续3-4分钟),但Dataproc的YARN监控显示CPU使用率最高仅约70%。

原因分析

  1. 统计口径差异

    • YARN仅统计自身管理的容器(Spark Executor的JVM进程)的CPU使用率,不包含节点上系统进程、Dataproc后台服务(如监控Agent、日志采集进程等)的CPU占用;而htop统计的是节点所有进程的总CPU使用率,当Executor满负荷时,加上后台服务的占用会让htop显示100%,但YARN只统计容器部分,数值自然更低。
    • 此外,YARN默认基于虚拟核计算使用率,若Worker节点开启超线程,YARN统计的是虚拟核的使用情况,而htop显示的是物理核的实际占用,二者维度不同也会造成数值偏差。
  2. 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%左右的使用率符合预期。
  3. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.19 12:37:20