Spark违反storageFraction阈值驱逐存储内存的异常问题
问题描述
我在Dataproc集群上运行小型PySpark实验,无法将动态内存分配逻辑与Spark UI观测结果对应。
集群配置
spark.executor.instances = 4 spark.executor.memory = 6G spark.executor.memoryOverHead = 2G spark.executor.pyspark.memory = 2G spark.memory.fraction = 0.6 spark.memory.offHeap.size = 2G spark.memory.storageFraction = 0.5
存储内存计算
我计算出每个executor的堆内存中不可驱逐部分为:
(6*(1024)-300)*0.6*0.5 1753.2 MBs per executor
4个executor总不可驱逐存储内存约为6.85GB。
实验操作
随后我读取9.5GB数据集并缓存:
df_sales = spark.read.option("format","parquet").option("header",True).option("inferSchema",True).load("gs://monsoon-credittech.appspot.com/spark_datasets/sales_parquet") from pyspark import StorageLevel df_sales = df_sales.persist(StorageLevel.MEMORY_AND_DISK) df_sales.count()
Spark UI存储页显示该数据集全量存入内存,无磁盘溢出。
异常现象
按逻辑,存储内存不应被驱逐至6.85GB以下,除非手动unpersist。但执行内存密集型join操作并持久化结果后,存储页显示:df_sales的内存数据被完全驱逐至磁盘,df_merged仅3.4GB在内存,总存储内存远低于阈值。未调用unpersist,为何Spark违反spark.memory.storageFraction阈值驱逐存储内存?
问题解析
这是因为你对spark.memory.storageFraction的作用存在误解,它不是不可驱逐存储内存的下限,而是动态内存池中分配给存储内存的比例上限,同时存储内存和执行内存之间支持动态互相借用。
具体细节如下:
- Spark的堆内存(
spark.executor.memory)会先预留300MB固定内存,剩余部分乘以spark.memory.fraction(此处为0.6)得到统一的动态内存池,这个池子同时服务于存储内存(缓存数据)和执行内存(shuffle、join等计算操作)。 spark.memory.storageFraction(此处为0.5)是指动态内存池中最多有50%的空间可被存储内存长期占用,但执行内存可以抢占这部分空间——只要存储内存中的数据是可驱逐的(比如你使用的MEMORY_AND_DISK级别,数据可写入磁盘)。- 当执行内存需求激增时(比如内存密集型join),Spark会优先驱逐存储内存中的可驱逐数据,哪怕这会让存储内存使用量低于
storageFraction对应的比例,直到执行内存获得足够空间,或者存储内存只剩不可驱逐的数据(比如MEMORY_ONLY且序列化的RDD、标记为不可驱逐的缓存)。
你的场景中,df_sales使用MEMORY_AND_DISK存储级别,属于可驱逐数据,当join操作需要大量执行内存时,Spark会将它的缓存数据全部驱逐到磁盘以满足执行内存需求,这完全符合Spark的内存管理规则,并未违反阈值设定。
另外需要注意,spark.executor.pyspark.memory是单独分配给Python进程的内存,不属于JVM堆内存的动态内存池,因此不会影响上述内存分配逻辑。
内容的提问来源于stack exchange,提问作者figs_and_nuts
相关产品推荐
相关产品推荐

