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

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(),看抽样比例是否符合预期,再排查全量转换的问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 16:01:18