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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 17:57:27