重写Process.run后未异步执行,为何导致线性而非并行处理?
为什么子类化Process.run会导致多进程池线性执行?
我对Process.run进行了子类化实现,代码如下:
import multiprocessing as mp import time DELAY = 2 class NewProcess(mp.get_context().Process): def run(self) -> None: # add new kwarg to item[4] slot old_que = self._args[0] new_que = mp.SimpleQueue() while not old_que.empty(): item = old_que.get() new_que.put( ( item[0], item[1], item[2], # Function item[3], # Arguments item[4] # Keyword arguments | { "message": "Hello world!", }, ) ) # Recreate args self._args = new_que, *self._args[1:] # Continue as normal super().run() def delay(*args, **kwargs): time.sleep(DELAY) return args, kwargs if __name__ == "__main__": context = mp.get_context() context.Process = NewProcess with context.Pool(2) as pool: responses = [] start = time.perf_counter() for _ in range(2): resp = pool.apply_async( func=delay, args=tuple(range(3)), kwds={}, ) responses.append(resp) for resp in responses: resp.wait() responses = [resp.get() for resp in responses] total = time.perf_counter() - start assert total - DELAY < 1e-2, f"Expected to take circa {DELAY}s, took {total}s" assert responses == ( expected := list( ( (0, 1, 2), { "message": "Hello world!" } ) ) ), f"{responses=}!={expected}"
我预期delay函数异步执行,总耗时约DELAY(2秒),但实际脚本执行失败,报错如下:
Traceback (most recent call last): File "/home/vahvero/Desktop/tmp.py", line 54, in <module> assert total - DELAY < 1e-2, f"Expected to take circa {DELAY}s, took {total}s" AssertionError: Expected to take circa 2s, took 4.003754430001209s
问题原因
- 进程池任务队列被独占:multiprocessing.Pool的多个工作进程共享同一个任务队列获取任务。你在
run方法里通过while not old_que.empty()把队列里的所有任务一次性取出并转移到新队列,这会导致第一个启动的工作进程抢走了所有任务,第二个进程无任务可执行,只能等第一个进程处理完所有任务才会结束。原本并行的任务变成了串行处理,总耗时自然翻倍。 - 空队列判断的不可靠性:
old_que.empty()在多进程场景下并非线程/进程安全的判断,主进程可能还在往队列里追加任务,但这里的核心矛盾还是任务被单个进程独占。
替代方案
不要通过修改Process的run方法破坏队列机制,推荐用包装目标函数的方式实现需求:
def wrap_delay(*args, **kwargs): kwargs["message"] = "Hello world!" return delay(*args, **kwargs)
之后在apply_async中使用func=wrap_delay即可,既不影响进程池的并行逻辑,又能给任务添加额外参数。
内容的提问来源于stack exchange,提问作者vahvero
相关产品推荐
相关产品推荐

