PySpark执行count正常,写入CSV时触发Java堆内存不足错误求助
Spark写入CSV时堆内存不足问题的原因分析与解决方案
问题场景
处理超大规模Parquet数据集时,执行groupBy("symbol").pivot("DDate").sum("value")后调用count()返回325073条记录,一切正常,但替换为写入CSV操作(无论是否使用repartition)时,立即触发Java heap space错误。当前配置:
spark.driver.memory 4g spark.executor.memory 12g
数据处理代码:
spark.read\ .format("parquet")\ .load('/mnt/data/share/parquets/quotes-*.parquet')\ .withColumn("symbol", SQL.regexp_extract('filename', r'\/([^\/]*).csv$', 1))\ .withColumn("DDate", col("Date").cast(DateType()))\ .select(["symbol", "DDate", "value"]) .filter(col("DDate") >= effective_date) .groupBy("symbol").pivot("DDate").sum("value") # .count() # 正常执行 .repartition(1).write.csv("/tmp/test.csv") # 触发OOM
错误日志核心片段:
23/04/05 22:30:27 ERROR Executor: Exception in task 12.0 in stage 10.0 (TID 662) java.lang.OutOfMemoryError: Java heap space
核心原因分析
Pivot操作导致单条记录过大
groupBy+Pivot会将每个symbol对应的所有DDate列聚合为一行,如果effective_date后的日期范围较广,单条Row会包含成百上千个字段。即使总记录数只有30多万,但单条Row的内存占用会远超预期,在shuffle、分区合并或写入阶段,大量大Row堆积会直接耗尽Executor堆内存。分区操作的误区
repartition(1)强制将所有数据集中到一个Executor任务中,所有大Row同时加载到内存,必然触发OOM;- 即使省略
repartition,pivot后的shuffle过程中,默认分区的单分区数据量仍可能因大Row存在而超出内存承载能力。
内存配置的实际可用空间不足
Spark Executor的12G内存并非全部用于堆存储,需预留部分给Off-Heap内存、用户代码、缓存及JVM本身,实际可用堆内存可能远低于12G,无法容纳大Row的批量处理。
解决方案
1. 优化Pivot逻辑,减少单条Row宽度
- 提前过滤不必要的日期:如果业务不需要
effective_date之后的所有日期,先通过filter缩小日期范围,减少pivot生成的列数; - 拆分Pivot任务:按
symbol的前缀或其他规则拆分数据集,分批次执行pivot和写入,避免单批次处理过多大Row。
2. 调整Spark内存与分区参数
- 增大Executor内存及内存预留:
spark.executor.memory=16g spark.executor.memoryOverhead=4g # 设为executor内存的20%-30% - 调整shuffle分区数:增大
spark.sql.shuffle.partitions(默认200),比如设为1000,让每个分区的大Row数量减少:spark.sql.shuffle.partitions=1000
3. 修改写入策略
- 避免
repartition(1):保留默认分区,或根据数据量设置合理的分区数(比如与Executor核心数匹配); - 使用
coalesce代替repartition(若需减少分区):coalesce不会触发全量shuffle,仅合并现有分区,降低内存压力:.coalesce(4).write.csv("/tmp/test.csv")
4. 启用自适应执行计划
开启Spark自适应执行,让系统自动调整分区和执行策略,适配大Row场景:
spark.sql.adaptive.enabled=true spark.sql.adaptive.shuffle.targetPostShuffleInputSize=64m # 调整目标分区大小
内容的提问来源于stack exchange,提问作者KIC
相关产品推荐
相关产品推荐

