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()
三、预期与实际差异
预期
- 数据优先持久化到内存,内存不足时再写入磁盘
- HDFS容量会因数据持久化的磁盘溢出而被占用
实际表现
- Spark UI显示数据先写入磁盘而非内存(对应Storage tab截图)
- 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
相关产品推荐
相关产品推荐

