从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
相关产品推荐
相关产品推荐

