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

PySpark监控指标不符预期:Spark UI与HDFS UI疑问排查

PySpark持久化与UI表现的认知误区解析

一、Session配置

pyspark --master yarn --num-executors 4 --executor-memory 6G --executor-cores 3 --conf spark.dynamicAllocation.enabled=false --conf spark.exector.memoryOverhead=2G --conf spark.memory.offHeap.size=2G --conf spark.pyspark.memory=2G

(注:配置存在拼写错误:spark.exector.memoryOverhead应为spark.executor.memoryOverhead)

二、执行代码

#Reading the same file twice
df_sales = spark.read.option("format","parquet").option("header",True).option("inferSchema",True).load("gs://monsoon-credittech.appspot.com/spark_datasets/sales_parquet")
df_sales_copy = spark.read.option("format","parquet").option("header",True).option("inferSchema",True).load("gs://monsoon-credittech.appspot.com/spark_datasets/sales_parquet")
#caching one
from pyspark import StorageLevel
df_sales = df_sales.persist(StorageLevel.MEMORY_AND_DISK)

#merging the two read files
df_merged = df_sales.join(df_sales_copy,df_sales.order_id==df_sales_copy.order_id,'inner')
df_merged = df_merged.persist(StorageLevel.MEMORY_AND_DISK)
#calling an action to trigger the transformations
df_merged.count()

三、预期与实际差异

预期

  1. 数据优先持久化到内存,内存不足时再写入磁盘
  2. HDFS容量会因数据持久化的磁盘溢出而被占用

实际表现

  1. Spark UI显示数据先写入磁盘而非内存(对应Storage tab截图)
  2. HDFS UI仅显示占用1.97GB,无明显溢出占用

四、认知误区梳理与实际行为解释

1. 关于MEMORY_AND_DISK的持久化顺序与UI显示

Spark的MEMORY_AND_DISK策略确实是优先尝试内存存储,内存不足时才溢出到磁盘,但你的观察存在两个偏差:

  • 持久化是懒执行的:persist仅标记DataFrame需要持久化,直到触发count()这类action才会实际执行。而join操作会触发shuffle,过程中会生成临时磁盘文件用于数据交换,这部分临时文件可能被误判为持久化的磁盘存储。
  • UI显示的是最终持久化状态:如果内存足够容纳数据,磁盘占用可能是shuffle临时文件而非持久化数据;如果内存不足,持久化数据会溢出到Executor本地磁盘,但UI会同时显示内存和磁盘的占用量。你看到的“先写入磁盘”大概率是shuffle临时文件的提前生成,而非持久化逻辑异常。

2. 关于持久化磁盘的存储位置

这是核心认知误区:Spark持久化到磁盘时,写入的是Executor节点的本地磁盘,而非HDFS。

  • HDFS是分布式文件系统,用于存储长期数据;Executor本地磁盘是节点的本地存储,仅用于临时存储shuffle中间数据、持久化溢出数据等,这些数据不会被计入HDFS的容量占用,所以HDFS UI看不到明显变化。
  • 你看到的HDFS占用1.97GB是原始Parquet文件的存储,和持久化操作无关。

3. 结合配置与数据量的进一步分析

你的Executor配置为4个节点,每个分配6G堆内存,默认情况下Spark会将堆内存的50%分配给存储内存(可通过spark.memory.fraction调整),即每个Executor的存储内存为3G,总存储内存为12G。

  • 原始Parquet文件是9GB,但Parquet是列式压缩格式,加载到内存后的实际数据量可能小于9GB,理论上内存足够容纳df_sales的持久化数据。但df_merged是join后的结果,inner join会导致数据量膨胀(如果order_id存在重复),此时总数据量可能超过存储内存上限,从而触发持久化到本地磁盘,但这部分不会影响HDFS。

内容的提问来源于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.08.09 08:40:22