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
相关产品推荐
相关产品推荐

