Spark分区大小超过Executor内存未触发OOM的原因问询
为什么Spark分区远大于Executor内存却未触发OOM?
你的推测完全正确,核心原因是Spark的执行模型并非一次性加载整个分区到内存,而是逐批次、逐行处理数据,再结合几个关键机制,让你避开了OOM:
逐行流式处理逻辑:Spark的转换操作(
withColumn、filter等)基于「惰性执行+逐条处理」模式。对于CSV这种行式存储,Spark会逐行读取解析,处理完一条就释放对应内存(或过滤后直接丢弃),不会把整个1000MiB的分区一次性塞进Executor内存。Web UI里读写行数一致也印证了这一点——处理的就是读取的每条数据,没有在内存中囤积全量分区数据。缓存未实际生效:你调用了
df_covid.cache(),但后续通过withColumn生成了新的DataFrame,原缓存的DF并没有被触发计算(count()是针对转换后的新DF)。就算缓存被触发,Spark默认的缓存策略是MEMORY_AND_DISK,内存不足时会自动把数据刷到磁盘,不会直接引发OOM。分区数据的分批处理:即使单个任务处理大分区,Spark内部也会把分区拆分成多个小批次(默认按1万条/批)来处理,单个批次的内存占用远小于整个分区的大小。你的Executor有2核,每个核处理一个批次,内存压力被进一步分散。
过滤操作的内存回收:
filter操作过滤掉了绝大多数数据,被过滤的记录会立即被垃圾回收(GC),不会长期占用内存。而最终的count()只需要统计行数,不需要保留所有处理后的结果,内存占用极低。
内容的提问来源于stack exchange,提问作者figs_and_nuts
相关产品推荐
相关产品推荐

