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

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的情况下流式遍历内部元素:

  1. 先安装依赖:pip install ijson
  2. 把解压后的流喂给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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 05:36:03