多进程处理大文本文件:如何保持输出与输入记录顺序一致?
大文本文件并行处理后保持输出顺序的解决方案
问题背景
有一个约300GiB的文本文件,包含可变行数的表头和数据记录。采用Python Multiprocessing做并行处理(CPU密集型任务),当前代码可正常运行,但输出记录顺序与输入不一致,需要高效实现顺序一致的方案,同时了解标准高效的实现架构。
示例输入文件 input.txt
# header # etc... (the number of lines in the header can vary) record #1 record #2 record #3 record #4 record #5 record #6 record #7 record #8 record #9 ...
示例代码
#!/usr/bin/env python import sys import multiprocessing def worker(queue,lock,process_id): while True: data = queue.get() if data == None: break # 模拟CPU密集型处理 for x in range(1000000): data.split() lock.acquire() print( data.rstrip() + " computed by process #" + process_id ) sys.stdout.flush() lock.release() if __name__ == '__main__': queue = multiprocessing.Queue() lock = multiprocessing.Lock() workers = [] num_processes = 4 for process_id in range(num_processes-1): p = multiprocessing.Process(target=worker, args=(queue,lock,str(process_id))) p.start() workers.append(p) with open('input.txt') as handler: try: # 读取表头 line = next(handler) while line.startswith('#'): line = next(handler) # 向队列发送记录 while True: queue.put(line) line = next(handler) except StopIteration: for p in workers: queue.put(None) finally: for p in workers: p.join()
当前输出
record #4 computed by process #2 record #3 computed by process #1 record #1 computed by process #0 record #2 computed by process #3 record #5 computed by process #2 record #6 computed by process #1 record #7 computed by process #0 record #8 computed by process #3 record #9 computed by process #2
一、实现输出顺序与输入一致的高效方案
核心思路是给每条记录绑定序号标记,Worker处理完成后返回(序号、结果),主进程按序号排序输出,或维护有序缓冲区批量输出。
方案1:使用multiprocessing.Pool.imap(轻量场景推荐)
Pool.imap会自动保留输入顺序,底层已封装结果排序逻辑,无需手动管理队列:
#!/usr/bin/env python import sys import multiprocessing def process_line(line): # 模拟CPU密集型处理 for x in range(1000000): line.split() return line.rstrip() + " computed by process #" + str(multiprocessing.current_process().pid % 4) if __name__ == '__main__': num_processes = 4 with open('input.txt') as handler: # 读取并输出表头 line = next(handler) while line.startswith('#'): print(line.rstrip()) line = next(handler) # 构造记录迭代器(从第一条非表头记录开始) def record_iterator(): yield line for l in handler: yield l # 用Pool.imap保持顺序处理 with multiprocessing.Pool(num_processes) as pool: for result in pool.imap(process_line, record_iterator()): print(result) sys.stdout.flush()
方案2:带序号的任务队列+结果排序(复杂任务场景)
手动管理进程时,给每条记录加序号,Worker返回(序号、结果),主进程收集后按序号排序输出:
#!/usr/bin/env python import sys import multiprocessing from queue import Empty def worker(task_queue, result_queue): while True: try: idx, data = task_queue.get(timeout=1) if idx is None: break except Empty: continue # 模拟CPU密集型处理 for x in range(1000000): data.split() result_queue.put( (idx, data.rstrip() + " computed by process #" + str(multiprocessing.current_process().pid % 4)) ) if __name__ == '__main__': num_processes = 4 task_queue = multiprocessing.Queue(maxsize=num_processes*2) result_queue = multiprocessing.Queue() workers = [] # 启动Worker进程 for _ in range(num_processes): p = multiprocessing.Process(target=worker, args=(task_queue, result_queue)) p.start() workers.append(p) with open('input.txt') as handler: # 读取并输出表头 line = next(handler) while line.startswith('#'): print(line.rstrip()) line = next(handler) # 发送带序号的任务 idx = 0 try: task_queue.put( (idx, line) ) idx +=1 for line in handler: task_queue.put( (idx, line) ) idx +=1 except StopIteration: pass # 发送终止信号 for _ in range(num_processes): task_queue.put( (None, None) ) # 收集结果并按序号排序 results = [] for _ in range(idx): results.append( result_queue.get() ) # 按序号排序后输出 for res in sorted(results, key=lambda x: x[0]): print(res[1]) sys.stdout.flush() # 等待Worker进程结束 for p in workers: p.join()
二、更标准高效的大文件并行处理架构
针对300GiB级别的超大型文件,可结合以下优化点提升效率:
- 分块读取替代逐行读取:按固定大小的块读取文件,Worker内部拆分行(处理跨行边界),减少进程间通信的开销。
- 避免全局锁与实时输出:让Worker将结果写入各自的临时文件,主进程最后按顺序合并临时文件,完全消除锁竞争。
- 使用
Pool的批量方法:Pool.map/Pool.starmap比手动管理队列更简洁,Python内部已做性能优化;若追求吞吐量,可搭配imap_unordered手动排序。 - 严格控制内存:保持迭代式处理,避免一次性加载过多数据到内存,确保内存占用稳定。
内容的提问来源于Stack Exchange,提问作者Fravadona
相关产品推荐
相关产品推荐

