Python multiprocessing结合Queue无打印、不写入文件问题求助
问题描述
尝试对for循环做并行化改造,将计算结果写入单文本文件,参考多进程场景下单文件队列写入、多进程安全写文件的技术方案编写实现代码,初始代码如下:
#multiprocessing version import multiprocessing as mp import time import numpy as np from itertools import product from functools import partial fn = './rules_res_3.txt' def worker(arg, q): '''模拟耗时计算逻辑''' w1, w2, w3= arg res=str(w1) +" " + str(w2) + " " + str(w3) q.put(res) return res def listener(q): '''监听队列消息,写入文件''' with open(fn, 'a') as f: while True: m = q.get() print(m) if m == 'kill': break f.write(m + '\n') # f.flush() def run(): r1,r2,r3 = np.arange(0.9, 1.5, 0.1), np.arange(0.9, 1.1, 0.1), np.arange(0, 1, 0.1) params=product(r1, r2, r3) with mp.Manager() as manager: pool = mp.Pool() # 默认根据CPU核心数创建进程池大小 message_queue = manager.Queue() # 用于向文件写入监听器传递消息的队列 pool.apply_async(listener, (message_queue, )) # 优先启动文件监听写入进程 pool.map(partial(worker, q=message_queue), params) # 拆分任务负载分发给worker进程 run()
运行异常
代码运行后出现两个问题:
listener函数内的print语句无任何控制台输出- 程序未向目标文件
./rules_res_3.txt写入任何内容,无法确认workers和listener进程是否正常运行
故障原因
代码存在三个核心问题:
- worker进程全部执行完成后,未向消息队列发送
kill终止信号,listener进程会一直阻塞在q.get()等待新消息,无法正常退出 - 未显式执行进程池关闭、等待任务完成的操作,
Manager上下文退出时会直接销毁队列、进程池等资源,listener进程还没来得及消费队列中的消息就被强制终止 - 文件写入默认走系统缓冲区缓存,没有正常走完文件关闭流程的情况下,缓冲区内容不会自动刷入磁盘,导致文件为空
修复后代码
import multiprocessing as mp import numpy as np from itertools import product from functools import partial fn = './rules_res_3.txt' def worker(arg, q): '''模拟耗时计算逻辑''' w1, w2, w3= arg res = f"{w1} {w2} {w3}" q.put(res) return res def listener(q): '''监听队列消息,写入文件''' with open(fn, 'a', encoding='utf-8') as f: while True: m = q.get() if m == 'kill': break print(m) f.write(m + '\n') f.flush() # 可选,实时刷入磁盘避免进程异常丢数据 def run(): # 首次运行先清空旧结果文件,避免历史内容干扰 with open(fn, 'w', encoding='utf-8') as f: pass r1,r2,r3 = np.arange(0.9, 1.5, 0.1), np.arange(0.9, 1.1, 0.1), np.arange(0, 1, 0.1) params = product(r1, r2, r3) with mp.Manager() as manager: pool = mp.Pool() message_queue = manager.Queue() # 启动监听写入进程 pool.apply_async(listener, (message_queue, )) # 分发计算任务 pool.map(partial(worker, q=message_queue), params) # 所有计算任务完成后,发送终止信号 message_queue.put('kill') # 关闭进程池,等待所有任务(包括listener)执行完成 pool.close() pool.join() if __name__ == '__main__': run()
跨平台注意:Windows系统下多进程代码必须放在
if __name__ == '__main__':判断下执行,否则会出现递归创建进程的报错,建议所有平台都保留该判断保证兼容性。
内容的提问来源于stack exchange,提问作者JP Zhang
相关产品推荐
相关产品推荐

