PySpark处理超宽Parquet文件遇Java堆内存溢出问题求助
问题
我有一个约10GB、包含25000列的大型Parquet文件,想要查看并将部分行转换为CSV格式。尝试过parquet-tools、fastparquet、pandas等工具均失败,改用PySpark后却遇到Java堆内存溢出错误(java.lang.OutOfMemoryError: Java heap space)。我的机器配有96GB内存,运行Python前已执行:
export JAVA_OPTS="-Xms36g -Xmx90g"
同时也尝试将driver memory设置为80GB,使用的代码如下:
from pyspark.sql import SQLContext from pyspark import SparkContext from pyspark.sql.types import * sc = SparkContext(appName="foo") sqlContext = SQLContext(sc) sc._conf.set('spark.driver.memory', '80g') readdf = sqlContext.read.parquet('dataset.parquet') readdf.head(2)
报错信息显示Executor的Java堆内存不足,同时有计划字符串过长的警告:
23/02/01 20:48:43 WARN package: Truncated the string representation of a plan since it was too large. This behavior can be adjusted by setting 'spark.sql.debug.maxToStringFields'. 23/02/01 20:48:48 ERROR Executor: Exception in task 7.0 in stage 3.0 (TID 13)20] java.lang.OutOfMemoryError: Java heap space ...(省略重复报错栈)
处理建议
1. 正确设置Spark内存配置(关键)
你在创建SparkContext之后设置spark.driver.memory是无效的——driver内存必须在Spark初始化前指定。改用SparkSession(推荐的新API)提前配置参数:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("parquet_process") \ .config("spark.driver.memory", "80g") \ # 分配给driver的内存 .config("spark.executor.memory", "16g") \ # 每个executor的内存,根据机器核数调整,比如16核机器给4个executor,每个4核16g .config("spark.executor.cores", "4") \ .config("spark.executor.memoryOverhead", "4g") \ # 额外的非堆内存开销,一般为executor内存的10%-25% .getOrCreate()
或者通过spark-submit命令行传递参数,优先级更高:
spark-submit --driver-memory 80g --executor-memory 16g --executor-cores 4 --executor-memory-overhead 4g your_script.py
2. 关闭Parquet矢量化读取(临时应急)
报错发生在矢量化读取器初始化阶段,超宽表(25000列)可能触发了矢量化读取的内存瓶颈,临时关闭该功能:
spark.conf.set("spark.sql.parquet.enableVectorizedReader", "false")
3. 只读取必要的列
25000列全量读取会极大消耗内存,如果只需要部分列,明确指定列名可以大幅降低内存占用:
# 示例:只读取需要的列名 target_cols = ["col_name_1", "col_name_2", "col_name_3"] # 替换为实际需要的列 readdf = spark.read.parquet('dataset.parquet').select(*target_cols) readdf.show(2) # 查看前2行
4. 调整元数据处理配置
针对计划字符串过长的警告,以及元数据读取内存溢出的问题,调整对应参数:
# 匹配列数设置字符串截断阈值 spark.conf.set("spark.sql.debug.maxToStringFields", "25000") # 增大元数据读取的上限 spark.conf.set("spark.sql.parquet.metadata.read.max", "1000000")
5. 分批处理并导出CSV
不要一次性处理全量数据,通过采样或限行数的方式导出部分数据:
- 采样少量数据查看:
# 随机采样0.1%的数据 sample_df = readdf.sample(fraction=0.001, seed=42) sample_df.write.csv("sample_output.csv", header=True, sep=",")
- 导出前N行:
# 导出前100行到CSV readdf.limit(100).write.csv("top_100_rows.csv", header=True, sep=",")
内容的提问来源于stack exchange,提问作者Vishaal
相关产品推荐
相关产品推荐

