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

Python如何拆分大型JSONL文件并分配给多进程执行数据筛选任务

大文件JSONL多进程分片处理实现方案

问题背景

你当前的实现采用主进程分发任务到队列的模式,对于14G的大文件来说,主进程单线程读文件、解析、推送任务的逻辑会成为整个流程的瓶颈,没有充分利用多核性能。按进程分配文件分片的逻辑可以避免这一问题,每个进程独立读取自己负责的文件片段,大幅降低进程间通信的开销。

现有代码可优化点

  • LINES_PER_PROCESS 没有做整数转换,会得到浮点数无法直接用做行数计数
  • 输出文件strokes_DE_file.json当前是只读模式打开,且写线程启动后with块立刻结束关闭文件,会触发写文件报错
  • 所有待处理数据都要经过in_q传递,序列化和IPC开销大,大文件场景下效率低

分片处理实现逻辑

核心思路

预先计算每个进程需要处理的行数范围,每个worker进程独立打开源文件,跳过自己负责范围之外的行,只处理对应分片内的数据,不需要主进程分发任务。

完整修改后代码

import multiprocessing as mp
import threading
import json
import reverse_geocoder as rg

# 预先计算单进程处理行数,取整
TOTAL_LINES = 12137928
WORKER_COUNT = mp.cpu_count() - 1
LINES_PER_PROCESS = TOTAL_LINES // WORKER_COUNT

def geocode_worker(file_path, start_line, end_line, out_q):
    """worker独立读取自己负责的分片范围"""
    rg.search(()) # 预先加载reverse_geocoder的数据集,避免子进程重复加载
    with open(file_path, 'r', encoding='utf-8') as f:
        # 跳过起始行之前的内容
        for _ in range(start_line):
            next(f)
        # 处理自己负责的行范围
        for line_num in range(start_line, end_line):
            try:
                line = next(f)
                data = json.loads(line)
                strokes = data['strokes']
                for strike in strokes:
                    strike_location = (strike['lat'], strike['lon'])
                    if rg.search(strike_location)[0]['cc'] == 'DE':
                        out_q.put(f"{json.dumps(strike)}\n")
            except StopIteration:
                break
    # 处理完成发送结束信号
    out_q.put(None)

def file_write_worker(out_q, output_path, worker_count):
    """独立写进程负责写结果,避免多进程写文件冲突"""
    with open(output_path, 'w', encoding='utf-8') as f:
        finish_cnt = 0
        while finish_cnt < worker_count:
            msg = out_q.get()
            if msg is None:
                finish_cnt += 1
                continue
            f.write(msg)

def get_germany_strokes(jsonl_file, output_file):
    out_q = mp.Queue()
    processes = []
    # 为每个worker分配行范围
    for worker_idx in range(WORKER_COUNT):
        start_line = worker_idx * LINES_PER_PROCESS
        # 最后一个进程处理剩下的所有行,避免整除问题丢数据
        end_line = TOTAL_LINES if worker_idx == WORKER_COUNT -1 else (worker_idx +1)*LINES_PER_PROCESS
        p = mp.Process(
            target=geocode_worker,
            args=(jsonl_file, start_line, end_line, out_q)
        )
        processes.append(p)
        p.start()
    
    # 启动写线程
    write_thread = threading.Thread(
        target=file_write_worker,
        args=(out_q, output_file, WORKER_COUNT)
    )
    write_thread.start()
    
    # 等待所有worker结束
    for p in processes:
        p.join()
    # 等待写线程结束
    write_thread.join()

if __name__ == "__main__":
    get_germany_strokes('data-2021-09-29.jsonl', 'strokes_DE_file.jsonl')

关键修改说明

  • 去掉了负责分发任务的in_q,改为每个worker独立读取自己负责的行范围,大幅降低进程间通信开销
  • 修复了输出文件打开模式错误的问题,写线程全程持有打开的文件对象直到所有处理完成
  • 最后一个worker的结束行数直接设为总行数,避免整数整除导致剩下的行没有被处理
  • 预先在子进程加载reverse_geocoder的数据集,避免每次调用rg.search都重复加载
    注:Windows系统下运行multiprocessing代码必须加if __name__ == "__main__"入口判断,否则会出现进程启动失败的问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 22:54:09