Spark如何处理大于Executor内存的分区?实践疑问解析
Spark 实验疑问解答
实验背景
- 数据集:磁盘大小12.3GB、内存反序列化大小3.6GB、序列化大小1964.9MB的CSV文件
- 集群配置:1个Executor,
spark.executor.memory=500mb,spark.executor.cores=2,spark.sql.files.maxPartitionBytes=13GB,关闭动态分配与自适应执行 - 操作:读取CSV得到2个分区(单分区磁盘约6GB,序列化后内存约880MB),执行
df.persist(StorageLevel.DISK_ONLY)和df.count(),观察到:任务未读完数据时无缓存、无数据溢出、存储页无数据
疑问1:Executor总内存仅500MB(300MB执行+存储内存+200MB用户内存),如何并行处理两个约880MB的分区?
Spark处理数据是流式逐块处理,并非把整个分区一次性加载到内存中:
- 执行
count()时,每个任务仅需逐行/逐数据块读取分区内容,完成计数后立即丢弃已处理的数据,不需要将880MB的整个分区数据长期驻留内存 - 配置2个Executor cores后会并行运行两个任务,但每个任务的内存占用仅取决于当前处理的数据块大小(远小于880MB),而非整个分区的总大小
- 加上你设置的持久化级别是
DISK_ONLY,Spark不会把处理中的数据暂存到内存存储区,进一步降低了内存压力
疑问2:读取的数据不在存储页、Executor内存中,也未溢出,数据存储在哪里?
数据存在Spark的磁盘临时存储目录,细节如下:
DISK_ONLY持久化级别明确要求Spark将处理完成的分区数据写入磁盘缓存,而非内存- 任务未完全读取数据时“无缓存”是正常行为:Spark采用边处理边持久化的逻辑,只有当整个分区的数据处理完成后,才会将完整的分区数据写入磁盘缓存目录(默认是系统临时目录,比如Linux下的
/tmp/spark-*路径) - 没有数据溢出是因为从一开始就没打算把数据放到内存里,所有持久化数据直接写入磁盘,不存在内存溢出的前提
内容的提问来源于stack exchange,提问作者figs_and_nuts
相关产品推荐
相关产品推荐

