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

PySpark在Dataproc读取GCS压缩CSV仅获部分记录,如何读取全部行?

问题描述

我有一个Google Dataproc作业,用于从Google Cloud Storage读取CSV文件,该文件的头部信息如下:

Content-type : application/octet-stream
Content-encoding : gzip
FileName: gs://test_bucket/sample.txt(文件无gz扩展名但已压缩)

以下代码运行成功,但DataFrame的记录数(9k)与文件实际记录数(100k)不匹配,似乎仅读取了前9k行:

self.spark :SparkSession= SparkSession.builder.appName("app_name"). \
                                                            config("spark.executor.memory","4g") \
                                                            .config("spark.hadoop.fs.gs.impl", "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem") \
                                                            .config("spark.hadoop.google.cloud.auth.service.account.enable", "true") \
                                                            .config("spark.hadoop.fs.gs.inputstream.support.gzip.encoding.enable", "true") \
                                                            .config("spark.sql.legacy.timeParserPolicy", "CORRECTED") \
                                                            .config("spark.driver.memory","4g").getOrCreate()

df = (self.spark.read.format("csv")
    .schema(schema)
    .option("mode", 'PERMISSIVE')  
    .option("encoding", "UTF-8")
    .option("columnNameOfCorruptRecord", '_corrupt_record')
    .load(self.file_path) )

print("df total count: ", df.count())
解决方法

1. 强制指定压缩格式

由于文件没有.gz扩展名,Spark可能无法自动识别gzip压缩格式,导致仅读取了部分解压后的数据。在读取CSV时显式指定压缩格式:

df = (self.spark.read.format("csv")
    .schema(schema)
    .option("mode", 'PERMISSIVE')  
    .option("encoding", "UTF-8")
    .option("columnNameOfCorruptRecord", '_corrupt_record')
    .option("compression", "gzip")  # 新增该行,强制使用gzip解码
    .load(self.file_path) )

2. 验证文件完整性

先确认GCS上的文件没有损坏,在终端执行以下命令,检查解压后的总行数:

gsutil cat gs://test_bucket/sample.txt | gunzip | wc -l

如果命令返回的行数不是100k,说明文件本身存在问题,需要重新上传完整的压缩文件。

3. 检查Schema匹配度

PERMISSIVE模式下,若Schema定义与CSV实际列不匹配,不兼容的行会被放入_corrupt_record列,这部分行不会被count()统计。执行以下代码查看脏数据量:

print("脏数据行数: ", df.filter("_corrupt_record is not null").count())

根据结果调整Schema,确保与CSV列完全匹配。

4. 调整Spark读取分区(可选)

如果文件体积过大,可调整分区参数确保并行读取完整数据:

# 设置单个分区最大字节数(示例为128MB,可根据文件大小调整)
self.spark.conf.set("spark.sql.files.maxPartitionBytes", "134217728")

df = (self.spark.read.format("csv")
    .schema(schema)
    .option("mode", 'PERMISSIVE')  
    .option("encoding", "UTF-8")
    .option("columnNameOfCorruptRecord", '_corrupt_record')
    .option("compression", "gzip")
    .load(self.file_path) )

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 20:54:51