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

Python3如何分块读取超内存x-bzip2 StreamingObject完整文件

问题原因
  • 你的parse方法每次执行都会重新创建bz2.BZ2File实例,且通过with上下文管理,方法执行结束后会自动关闭解压缩流;而S3返回的StreamingBody是仅支持单次消费的IO流,读取指针走到末尾后无法重置,第二次调用parse自然读不到任何数据
  • 现有逻辑没有维持解压缩流的全局状态,每次读取都相当于重头解码,但流已经被消费过,自然只能拿到首个分块的数据
修复后的实现

直接把解压缩流的初始化放到类的构造方法里,做成可迭代对象,每次迭代返回一个分块的解析结果:

import bz2
import json

class ParseStreamingZip(object):
    def __init__(self, obj, chunk_size=(1024*1024*5), dec='utf8'):
        "Chunksize: 默认5MB,dec为解码格式默认utf8"
        self.meta = obj
        self.chunk_size = chunk_size
        self.dec = dec
        self.streaming_obj = obj['Body']
        self.streaming_obj.set_socket_timeout(9999999)
        # 初始化解压缩流,不要放到with里自动关闭,等全量读完再关
        self.bz2_f = bz2.BZ2File(self.streaming_obj, 'rb')

    def __iter__(self):
        return self

    def __next__(self):
        # 读取当前分块的所有行
        content = self.bz2_f.readlines(self.chunk_size)
        # 读到空代表文件结束,停止迭代
        if not content:
            self.bz2_f.close()
            raise StopIteration
        # 解析当前分块的json行
        output = []
        for line in content:
            try:
                jline = json.loads(line.decode(self.dec).strip('\n'))
                output.append(jline)
            except Exception as e:
                print('Caught: ', e)
                pass
        return output

    # 兼容手动调用parse的方式
    def parse(self):
        return next(self, [])
全量分块处理调用示例
import pandas as pd

filename = 'wls_day-78.bz2'    

# 获取S3对象
response_obj = s3.get_object(Bucket=dataset_metadata['bucket'], Key=filename)
print('+ object received.')

# 打印对象信息
content_type = response_obj['ResponseMetadata']['HTTPHeaders']['content-type']
content_length = int(response_obj['ResponseMetadata']['HTTPHeaders']['content-length'])
print('content type & length:', content_type, content_length, "({:.1f} MB)".format(int(content_length)/1024/1024))

# 初始化流式解析类
parser = ParseStreamingZip(response_obj, chunk_size=1024*1024*5)

# 循环处理所有分块
all_dfs = []
for chunk_idx, parsed_chunk in enumerate(parser):
    df_chunk = pd.DataFrame(parsed_chunk)
    print(f"处理第{chunk_idx+1}个分块,内存占用:{df_chunk.memory_usage().sum()/1024/1024:.2f}MB")
    # 这里可以直接对df_chunk做增量处理,比如写入数据库、写入本地文件,不需要全量存内存
    all_dfs.append(df_chunk)

# 如果需要合并所有分块为全量DataFrame(内存足够的情况下)
full_df = pd.concat(all_dfs, ignore_index=True)
注意事项
  • 如果文件过大不需要全量存入内存,可以直接在循环里对每个分块做增量处理,不用把所有df_chunk都存到all_dfs列表里
  • 如果遇到网络中断需要重读,需要重新调用s3.get_object获取新的StreamingBody再初始化解析类
  • 解码异常的行可以根据需求做落盘留存,避免数据丢失

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 08:06:03