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

从S3分块读取20GB Gzip压缩JSON文件转Parquet时遇EOFError问题排查

解决S3大Gzip压缩JSON文件转Parquet时的EOFError问题

首先,你的原始代码出现EOFError的核心原因很明确:gzip是针对整个文件的流式压缩格式,你用S3的iter_chunks()获取的是原始压缩文件的二进制分块,这些分块并不是独立的gzip压缩单元。直接对单个不完整的压缩块调用gzip.decompress(),自然会因为找不到完整的流结束标记而报错。

方案一:流式解压+批量处理(无本地文件)

我们可以用gzip.GzipFile直接包装S3的响应流,实现流式解压,同时逐行读取JSON数据,累积到一定批量后再转成Parquet保存到S3,这样既不会占用过多内存,也能避免分块解压的错误。

示例代码:

import json
import pandas as pd
import gzip
from io import TextIOWrapper
from your_s3_module import S3  # 替换成你的S3连接模块

BATCH_SIZE = 100000  # 根据内存情况调整批量大小
input_bucket = "your-input-bucket"
object_key = "path/to/large-file.json.gz"
output_bucket = "your-output-bucket"

with S3() as s3:
    obj = s3.s3_client.get_object(Bucket=input_bucket, Key=object_key)
    # 用GzipFile包装S3的二进制流,再转成文本流
    with gzip.GzipFile(fileobj=obj['Body'], mode='rb') as gz_file:
        text_stream = TextIOWrapper(gz_file, encoding='utf-8')
        
        batch_data = []
        count = 0
        
        for line in text_stream:
            line = line.strip()
            if not line:
                continue
            try:
                json_obj = json.loads(line)
                batch_data.append(json_obj)
            except json.JSONDecodeError as e:
                print(f"跳过无效JSON行: {line}, 错误: {e}")
                continue
            
            # 当批量达到设定大小,转成Parquet保存
            if len(batch_data) >= BATCH_SIZE:
                count += 1
                df = pd.DataFrame(batch_data)
                output_key = f"parquet_test/df_{count}.parquet.gzip"
                s3_url = f"s3://{output_bucket}/{output_key}"
                df.to_parquet(s3_url, compression='gzip', index=False)
                batch_data = []  # 重置批量
        
        # 处理最后一批剩余数据
        if batch_data:
            count += 1
            df = pd.DataFrame(batch_data)
            output_key = f"parquet_test/df_{count}.parquet.gzip"
            s3_url = f"s3://{output_bucket}/{output_key}"
            df.to_parquet(s3_url, compression='gzip', index=False)

这个方案的优势:

  • 完全流式处理,不需要加载整个文件到内存
  • 避免了分块解压的错误,因为gzip.GzipFile会处理整个压缩流
  • 可以灵活调整BATCH_SIZE来平衡内存占用和处理效率

方案二:优化你的Pandas方案(避免本地文件)

你尝试的pd.read_json方法其实是可行的,但可以改进来避免生成本地文件——直接用Pandas的to_parquet写入S3,不需要先存本地。这里需要确保你的环境已经安装了s3fs(pip install s3fs),Pandas可以通过s3fs直接读写S3路径。

优化后的代码:

import pandas as pd
from your_s3_module import S3Connect  # 替换成你的S3连接模块

CHUNKSIZE = 500000  # 根据内存调整,不用设太大,50万行一般足够
input_s3_path = "s3://your-input-bucket/path/to/large-file.json.gz"
output_bucket = "your-output-bucket"

count = 0
# 直接读取S3上的压缩JSON文件,分块处理
for df in pd.read_json(input_s3_path, lines=True, chunksize=CHUNKSIZE):
    count += 1
    output_key = f"parquet_test/df_{count}.parquet.gzip"
    output_s3_path = f"s3://{output_bucket}/{output_key}"
    df.to_parquet(output_s3_path, compression='gzip', index=False)

这个方案更简洁,Pandas内部已经处理了gzip解压和流式读取的逻辑,不需要自己手动处理流。需要注意的点:

  • 确保安装了s3fs和pyarrow/fastparquet(Parquet引擎)
  • chunksize不要设置过大,否则会占用过多内存导致OOM

为什么原始代码会出错?

再强调一下:obj['Body'].iter_chunks()返回的是原始压缩文件的二进制片段,每个片段只是整个gzip流的一部分,没有完整的gzip头部和尾部信息。gzip.decompress()只能处理完整的gzip压缩包,所以单独解压这些片段必然会抛出EOFError。

内容的提问来源于stack exchange,提问作者M. Phys

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 09:29:05