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

PySpark拆分大数据帧为指定大小CSV时出现Py4JJavaError

问题

我尝试将包含700万条记录的PySpark DataFrame拆分为每个含80万条记录的CSV文件,已编写完整代码(如下),但执行时抛出Py4JJavaError,底层错误为java.lang.OutOfMemoryError: Java heap space(堆内存不足):

from pyspark.sql import SparkSession

options = {
    "pathGlobalFilter": "*.csv",
    "header": "True",
}

spark = SparkSession.builder.config("spark.driver.host","localhost").appName("CSV Reader").getOrCreate()

csv_path = "C:\\Users\\rajat.kapoor\\Desktop\\155_Lacs_Raw_Data\\OutputFiles_CSV"

RawData_Combined_Revolt_Only = spark.read.format("csv").options(**options).load(csv_path)

# Import necessary libraries
from pyspark.sql.functions import monotonically_increasing_id

# Define the chunk size
chunk_size = 800000

# Add a unique ID column to the dataframe
RawData_Combined_Revolt_Only = RawData_Combined_Revolt_Only.withColumn("id", monotonically_increasing_id())

# Repartition the dataframe based on the chunk size
RawData_Combined_Revolt_Only = RawData_Combined_Revolt_Only.repartition((RawData_Combined_Revolt_Only.count() / chunk_size) + 1)

# Write each partition to a separate CSV file
RawData_Combined_Revolt_Only.write.csv("C:\\Users\\rajat.kapoor\\Desktop\\Output PySpark Folder", header=True, mode="overwrite")
排查与解决方案

1. 核心问题定位

代码中直接调用RawData_Combined_Revolt_Only.count()会触发全量数据扫描,将计算结果返回到Driver端,700万条记录的统计操作会瞬间占用大量Driver内存,直接引发堆内存溢出,进而抛出Py4JJavaError。

2. 针对性优化方案

方案一:用ID分段替代提前count(推荐)

不需要提前计算总条数,利用新增的id字段进行逻辑分段,确保每个分区数据量接近目标chunk_size,同时避免Driver端的内存压力:

from pyspark.sql import SparkSession
from pyspark.sql.functions import monotonically_increasing_id, floor

options = {
    "pathGlobalFilter": "*.csv",
    "header": "True",
}

# 先临时调大Driver内存应急,根据机器实际配置调整
spark = SparkSession.builder.config("spark.driver.host","localhost")\
    .config("spark.driver.memory", "8g")\
    .appName("CSV Reader").getOrCreate()

csv_path = "C:\\Users\\rajat.kapoor\\Desktop\\155_Lacs_Raw_Data\\OutputFiles_CSV"
RawData_Combined_Revolt_Only = spark.read.format("csv").options(**options).load(csv_path)

chunk_size = 800000

# 新增唯一ID并计算分段键
df_with_id = RawData_Combined_Revolt_Only.withColumn("id", monotonically_increasing_id())
df_with_part = df_with_id.withColumn("partition_key", floor(df_with_id["id"] / chunk_size))

# 按分段键 repartition,每个分区对应一个CSV文件,最后删除临时字段
df_with_part.repartition("partition_key")\
    .drop("id", "partition_key")\
    .write.csv("C:\\Users\\rajat.kapoor\\Desktop\\Output PySpark Folder", header=True, mode="overwrite")

方案二:调大Spark内存参数

如果要保留原逻辑,必须调整Driver和Executor的内存配置,给Spark足够的运算空间:

spark = SparkSession.builder.config("spark.driver.host","localhost")\
    .config("spark.driver.memory", "16g")\ # 根据机器物理内存调整,比如16g/32g
    .config("spark.executor.memory", "16g")\
    .config("spark.executor.cores", "4")\
    .appName("CSV Reader").getOrCreate()

注意:机器的实际物理内存必须能支撑配置的内存值,否则会引发系统级OOM。

方案三:用coalesce替代repartition(适用于分区数减少场景)

如果源数据的分区数大于目标分区数,使用coalesce可以避免全量数据shuffle,大幅降低内存消耗:

# 粗略估算目标分区数,无需精确count
target_partitions = (7000000 // chunk_size) + 1
df_with_id.coalesce(target_partitions)\
    .write.csv("C:\\Users\\rajat.kapoor\\Desktop\\Output PySpark Folder", header=True, mode="overwrite")

3. 额外注意事项

  • 避免在Driver端执行大规模数据计算,尽量将逻辑下推到Executor分布式执行。
  • 如果源CSV文件本身已拆分为多个文件,可以直接复用原分区,无需额外repartition,减少shuffle开销。

内容的提问来源于stack exchange,提问作者RajatK350

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 05:55:39