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
相关产品推荐
相关产品推荐

