PySpark DataFrame写入CSV后丢失近半数行,求解决方案
Spark DataFrame行数缺失排查与解决
核心判断
两种操作(写入CSV、转Pandas)都出现相同的行数缺失,说明问题不在写入/转换环节,而是原始Spark DataFrame的实际行数就不足100万,或者数据生成/加载时存在异常。
排查与解决步骤
先确认DataFrame真实行数
不要依赖创建逻辑的预设值,直接用count()统计实际行数:print(df.count())如果输出不是100万,问题出在数据生成环节,直接回溯生成代码的逻辑。
检查数据生成逻辑
如果是自定义生成数据(比如随机生成、从数据源加载),排查以下点:- 是否有隐式的过滤、去重逻辑?
- 随机生成时是否存在重复值被自动去重?
- 循环或生成函数是否有逻辑错误,导致实际生成行数不足?
排查分布式分区异常
Spark是分布式计算框架,可能存在分区数据丢失的情况:# 查看分区总数 print(df.rdd.getNumPartitions()) # 查看每个分区的行数 df.rdd.mapPartitions(lambda iter: [len(list(iter))]).collect()如果某个分区行数为0或异常少,说明数据生成时的分区处理逻辑有问题,比如分区不均、部分分区生成失败。
验证CSV写入的完整性
Spark写入CSV时,默认会按分区生成多个part-*.csv文件,不要只统计单个文件的行数,需要把输出目录PresSample2.csv下所有CSV文件的行数相加,才是总数据量。转Pandas时的内存限制
如果转Pandas时Driver内存不足,可能会静默丢失数据(部分环境下不会报错):- 调大Driver内存参数,比如启动Spark时设置
--driver-memory 8g - 先抽样转换验证:
df.sample(fraction=0.1).toPandas(),看抽样比例是否符合预期,再排查全量转换的问题。
- 调大Driver内存参数,比如启动Spark时设置
内容的提问来源于stack exchange,提问作者Parisa Ebrahimifar
相关产品推荐
相关产品推荐

