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

