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

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进程是否正常运行
故障原因

代码存在三个核心问题:

  1. worker进程全部执行完成后,未向消息队列发送kill终止信号,listener进程会一直阻塞在q.get()等待新消息,无法正常退出
  2. 未显式执行进程池关闭、等待任务完成的操作,Manager上下文退出时会直接销毁队列、进程池等资源,listener进程还没来得及消费队列中的消息就被强制终止
  3. 文件写入默认走系统缓冲区缓存,没有正常走完文件关闭流程的情况下,缓冲区内容不会自动刷入磁盘,导致文件为空
修复后代码
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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 10:54:21