multiprocessing Manager共享进程池调用apply_async失败问题排查
我运行两个守护进程:一个向共享multiprocessing.Queue写入数据,另一个从该Queue读取数据并提交至multiprocessing.pool.apply_async,异步结果发送至另一个multiprocessing.Queue。
示例代码如下:
import multiprocessing from multiprocessing import Process import random import time def my_method(i): return i*i class DataFeeder: input_queue = None @staticmethod def stream_to_queue(iq): if DataFeeder.input_queue is None: DataFeeder.input_queue = iq while True: time.sleep(1) dat = random.choice(range(0, 50)) print(f"feeding {dat}") DataFeeder.input_queue.put(dat) class DataEater: input_queue = None results_queue = None pool = None @staticmethod def eat(iq, rq, p): if DataEater.input_queue is None: DataEater.input_queue = iq if DataEater.results_queue is None: DataEater.results_queue = rq if DataEater.pool is None: DataEater.pool = p while True: time.sleep(0.1) dat = DataEater.input_queue.get(0.1) # 100ms timeout print(f"eating {dat}") async_result = DataEater.pool.apply_async(my_method, (dat,)) print(f"async_result {async_result}") DataEater.results_queue.put_nowait(async_result) if __name__ == '__main__': with multiprocessing.Manager() as m: input_q = m.Queue() output_q = m.Queue() pool = m.Pool(8) dfp = Process(target=DataFeeder.stream_to_queue, args=(input_q,), daemon=True) dfp.start() dep = Process(target=DataEater.eat, args=(input_q, output_q, pool), daemon=True) dep.start() dep.join() dfp.join()
运行后出现如下堆栈跟踪:
Process Process-3: Traceback (most recent call last): File "/opt/anaconda3/envs/integration_tests/lib/python3.9/multiprocessing/process.py", line 315, in _bootstrap self.run() File "/opt/anaconda3/envs/integration_tests/lib/python3.9/multiprocessing/process.py", line 108, in run self._target(*self._args, **self._kwargs) File "/Users/chhunb/PycharmProjects/micromanager_server/umanager_server/multiproctest.py", line 54, in eat async_result = DataEater.pool.apply_async(my_method, (dat,)) File "<string>", line 2, in apply_async File "/opt/anaconda3/envs/integration_tests/lib/python3.9/multiprocessing/managers.py", line 816, in _callmethod proxytype = self._manager._registry[token.typeid][-1] AttributeError: 'NoneType' object has no attribute '_registry'
我尝试在eat()方法内本地创建进程池,但因async_result需通过results_queue与其他进程共享而无法实现;也尝试硬编码输入忽略DataFeeder,仍出现相同报错。该问题与CPython开源问题#80100高度相似,目前该问题仍处于开放状态,想知道后续Python版本是否修复此问题,是否需要更换设计模式来实现该功能?
问题原因
你遇到的错误是因为通过Manager创建的Pool代理不能跨进程传递。Manager生成的Pool是一个远程代理对象,它依赖于Manager的内部通信机制,当你把这个代理传递给另一个Process时,代理对象的_manager属性会变成None,导致调用apply_async时无法找到注册表,触发AttributeError。
CPython的#80100问题确实是关于跨进程传递Manager创建的Pool代理的缺陷,截至目前(包括Python 3.12)该问题仍未完全修复,官方暂无明确修复时间表。
解决方案:更换设计模式
不需要放弃异步处理,只需调整架构,避免跨进程传递Pool代理,同时解决AsyncResult共享的问题:
方案1:主进程管理结果收集,子进程只传递任务参数
修改逻辑,让DataEater只把需要处理的任务参数放入结果队列,由主进程负责提交到Pool并处理异步结果:
import multiprocessing from multiprocessing import Process import random import time def my_method(i): return i*i class DataFeeder: @staticmethod def stream_to_queue(iq): while True: time.sleep(1) dat = random.choice(range(0, 50)) print(f"feeding {dat}") iq.put(dat) class DataEater: @staticmethod def eat(iq, task_queue): while True: time.sleep(0.1) try: dat = iq.get(timeout=0.1) print(f"eating {dat}") task_queue.put(dat) except multiprocessing.queues.Empty: continue if __name__ == '__main__': with multiprocessing.Manager() as m: input_q = m.Queue() task_q = m.Queue() pool = multiprocessing.Pool(8) # 主进程本地创建Pool,不通过Manager dfp = Process(target=DataFeeder.stream_to_queue, args=(input_q,), daemon=True) dfp.start() dep = Process(target=DataEater.eat, args=(input_q, task_q), daemon=True) dep.start() # 主进程负责提交任务和处理结果 while True: try: dat = task_q.get(timeout=0.1) async_result = pool.apply_async(my_method, (dat,), callback=lambda res: print(f"result: {res}")) except multiprocessing.queues.Empty: continue
方案2:在DataEater进程内创建Pool,用共享队列传递计算结果而非AsyncResult
如果必须在子进程内处理任务提交,直接在eat方法内创建本地Pool,然后把计算完成的结果放入结果队列,而非传递AsyncResult(因为AsyncResult是本地对象,无法跨进程共享):
import multiprocessing from multiprocessing import Process import random import time def my_method(i): return i*i class DataFeeder: @staticmethod def stream_to_queue(iq): while True: time.sleep(1) dat = random.choice(range(0, 50)) print(f"feeding {dat}") iq.put(dat) class DataEater: @staticmethod def eat(iq, results_queue): # 子进程内本地创建Pool pool = multiprocessing.Pool(8) results = [] def collect_result(res): results_queue.put(res) print(f"result: {res}") while True: time.sleep(0.1) try: dat = iq.get(timeout=0.1) print(f"eating {dat}") # 提交任务时指定回调函数,直接把结果放入共享队列 pool.apply_async(my_method, (dat,), callback=collect_result) except multiprocessing.queues.Empty: continue if __name__ == '__main__': with multiprocessing.Manager() as m: input_q = m.Queue() output_q = m.Queue() dfp = Process(target=DataFeeder.stream_to_queue, args=(input_q,), daemon=True) dfp.start() dep = Process(target=DataEater.eat, args=(input_q, output_q), daemon=True) dep.start() # 主进程可以从output_q获取结果做后续处理 while True: try: res = output_q.get(timeout=0.1) # 处理结果逻辑 except multiprocessing.queues.Empty: continue
关键说明
- 避免跨进程传递Manager生成的Pool代理,这是当前Python multiprocessing的已知限制。
AsyncResult是进程本地对象,不能通过共享队列传递,要么在创建它的进程内处理结果,要么直接传递计算后的最终值。
内容的提问来源于stack exchange,提问作者bryantC

