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

如何用Python多进程并行解析gzip格式的JSON行文件?

并行解析gzip压缩的JSON行文件

实现思路

采用生产者-消费者模式:

  • 生产者进程负责读取gzip文件,将内容按固定行数分块,存入进程安全的队列
  • 4个消费者进程从队列中取出块,并行解析每行的JSON字符串
  • 通过队列传递结束信号,确保所有进程正常退出

这种方式既保证了任务大致平均分配(按块大小控制),又避免了逐行分发带来的进程间通信开销。

完整代码

import gzip
import json
from multiprocessing import Process, Queue

def consumer(queue):
    """消费者进程:从队列取块并解析JSON"""
    while True:
        batch = queue.get()
        # 收到结束信号则退出
        if batch is None:
            break
        for line in batch:
            try:
                # 解析JSON,这里可替换为你的业务逻辑
                data = json.loads(line.strip())
                # 示例:打印解析结果(实际场景可改为存储、计算等操作)
                # print(data)
            except json.JSONDecodeError as e:
                print(f"解析JSON失败: {e},行内容: {line[:50]}...")

def producer(file_path, queue, batch_size, num_consumers):
    """生产者进程:读取gzip文件并分块存入队列"""
    with gzip.open(file_path, 'rt', encoding='utf-8') as f:
        batch = []
        for line in f:
            batch.append(line)
            # 达到块大小则存入队列
            if len(batch) == batch_size:
                queue.put(batch)
                batch = []
        # 处理剩余不足一块的内容
        if batch:
            queue.put(batch)
    # 给每个消费者发送结束信号
    for _ in range(num_consumers):
        queue.put(None)

if __name__ == "__main__":
    # 配置参数
    FILE_PATH = "your_data.json.gz"  # 替换为你的gzip文件路径
    BATCH_SIZE = 250000  # 每个进程处理约25万行,对应100万行4个进程
    NUM_CONSUMERS = 4

    # 创建进程安全队列
    task_queue = Queue(maxsize=NUM_CONSUMERS * 2)  # 限制队列大小,防止内存溢出

    # 启动消费者进程
    consumers = []
    for _ in range(NUM_CONSUMERS):
        p = Process(target=consumer, args=(task_queue,))
        p.start()
        consumers.append(p)

    # 启动生产者进程
    producer_process = Process(target=producer, args=(FILE_PATH, task_queue, BATCH_SIZE, NUM_CONSUMERS))
    producer_process.start()

    # 等待生产者完成
    producer_process.join()

    # 等待所有消费者完成
    for p in consumers:
        p.join()

    print("所有JSON解析完成")

关键说明

  1. 分块策略:BATCH_SIZE设为250000,确保100万行大致平均分配给4个进程,剩余不足一块的内容会被最后一个处理的进程接收
  2. 队列控制:设置队列maxsize避免生产者过快写入导致内存占用过高
  3. 错误处理:加入JSONDecodeError捕获,防止单行解析失败导致进程崩溃
  4. 结束信号:生产者完成后向队列放入与消费者数量相等的None,确保每个消费者都能收到退出信号

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 01:57:38