You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.02 07:45:29