Python流式下载解压gz文件并将分块数据转JSON写入MongoDB方案问询
解决方案
核心逻辑是通过缓冲区拼接截断的分块内容,再根据解压后的数据格式选择对应解析方案,全程无需加载1.2GB全量数据到内存,内存占用可控制在几十MB以内。
场景1:解压后为JSON Lines格式(每行一个独立JSON对象,最常见的大数据传输格式)
绝大多数公开压缩JSON数据集都采用这种格式,处理成本最低:
- 维护一个字节缓冲区,每次将解压得到的分块内容追加到缓冲区
- 按换行符拆分缓冲区,完整行单独取出做JSON解析
- 拆分后剩余的不完整内容留在缓冲区,和下一次解压的内容拼接
- 所有分块处理完成后,处理缓冲区剩余的最后一段内容
代码示例:
import requests import json import zlib from pymongo import MongoClient # 初始化MongoDB连接 client = MongoClient('mongodb://localhost:27017/') db = client['你的数据库名'] col = db['你的集合名'] # 批量插入缓存,减少MongoDB IO次数,可根据实际情况调整大小 batch = [] BATCH_SIZE = 500 url = "https://something" d = zlib.decompressobj(zlib.MAX_WBITS | 16) # 初始化字节缓冲区 buffer = b'' with requests.get(url, stream=True, timeout=30) as r: r.raise_for_status() # 建议调大chunk_size到1M,128字节太小会大幅降低处理效率 for chunk in r.iter_content(chunk_size=1024*1024): # 解压当前分块 data = d.decompress(chunk) buffer += data # 按换行符拆分完整行 while b'\n' in buffer: line, buffer = buffer.split(b'\n', 1) # 跳过空行 line = line.strip() if not line: continue try: # 转成JSON字典 item = json.loads(line.decode('utf-8')) batch.append(item) # 攒够批量大小就插入MongoDB if len(batch) >= BATCH_SIZE: col.insert_many(batch) batch = [] except json.JSONDecodeError: # 格式错误的行可根据需求打日志或直接跳过 pass # 处理zlib剩余的解压数据 remaining = d.flush() if remaining: buffer += remaining # 处理缓冲区最后剩余的内容 if buffer.strip(): try: item = json.loads(buffer.strip().decode('utf-8')) batch.append(item) except json.JSONDecodeError: pass # 插入最后剩余的批次 if batch: col.insert_many(batch)
场景2:解压后是单一大JSON结构(比如一整个JSON数组)
这种情况需要用增量JSON解析库处理,推荐用ijson,可以在不加载完整JSON的情况下流式遍历内部元素:
- 先安装依赖:
pip install ijson - 把解压后的流喂给ijson做增量解析即可
代码示例:
import requests import zlib import ijson from pymongo import MongoClient client = MongoClient('mongodb://localhost:27017/') db = client['你的数据库名'] col = db['你的集合名'] batch = [] BATCH_SIZE = 500 url = "https://something" d = zlib.decompressobj(zlib.MAX_WBITS | 16) def decompress_stream(response): for chunk in response.iter_content(chunk_size=1024*1024): yield d.decompress(chunk) yield d.flush() with requests.get(url, stream=True, timeout=30) as r: r.raise_for_status() # 这里的'item'是JSON数组下的元素路径,根据你实际的JSON结构调整 for item in ijson.items(decompress_stream(r), 'item'): batch.append(item) if len(batch) >= BATCH_SIZE: col.insert_many(batch) batch = [] if batch: col.insert_many(batch)
如果你的数据是CSV等其他结构化格式,使用对应流式解析库(如pandas的read_csv分块读取)转成字典后,按同样的批量插入逻辑写入MongoDB即可。
内容的提问来源于stack exchange,提问作者Shahriar Tasnim
相关产品推荐
相关产品推荐

