PySpark on GCP Dataproc读取Gzip压缩GCS文件数据不全问题求助
问题原因分析
你遇到的问题核心在于流处理冲突与分片读取限制:
- 当启用
core:fs.gs.inputstream.support.gzip.encoding.enable=true时,GCS Hadoop连接器会在底层对带有Content-Encoding: gzip的对象自动解压,返回原始JSON流。 - 但Spark的JSON数据源依赖分片(split)实现并行计算,而gzip解压后的流属于不可随机访问的顺序流(无法seek到任意位置)。这导致Spark只能读取每个流的第一个分片内容,后续分片无法正确加载,最终仅处理了部分数据。
- 压缩级别改变时,压缩后的文件大小变化,对应Spark分片大小对应的解压后数据量也会改变,因此处理的内容占比随之波动;而平均值一致是因为已处理的分片内数据是完整且正确的。
解决方案
需要让Spark直接处理压缩文件,而非依赖Hadoop连接器的底层解压,具体修改如下:
1. 调整Dataproc工作流模板
移除core:fs.gs.inputstream.support.gzip.encoding.enable=true配置项,停止Hadoop连接器的自动解压行为:
gcloud dataproc workflow-templates set-managed-cluster $WORKFLOW_NAME \ # 保留其他原有配置,移除上述gzip相关的properties参数
2. 修改PySpark读取代码
在Spark读取JSON的操作中,显式指定压缩格式为gzip,让Spark数据源自行处理解压和并行分片逻辑:
decodedDF = spark.read.schema(globalSchema) \ .option("compression", "gzip") \ .json(<list-of-GCS-folders-containing-json-files>)
注意:
header=True参数对JSON数据源无效,可直接移除,避免不必要的配置干扰。
修改后,Spark会正确识别每个gzip压缩的JSON文件,利用自身的压缩处理逻辑完成解压,同时支持并行分片读取,确保所有数据都被完整处理。
内容的提问来源于stack exchange,提问作者Stéphane Janel
相关产品推荐
相关产品推荐

