如何用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解析完成")
关键说明
- 分块策略:
BATCH_SIZE设为250000,确保100万行大致平均分配给4个进程,剩余不足一块的内容会被最后一个处理的进程接收 - 队列控制:设置队列
maxsize避免生产者过快写入导致内存占用过高 - 错误处理:加入
JSONDecodeError捕获,防止单行解析失败导致进程崩溃 - 结束信号:生产者完成后向队列放入与消费者数量相等的
None,确保每个消费者都能收到退出信号
内容的提问来源于stack exchange,提问作者Felipe Lopes
相关产品推荐
相关产品推荐

