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

重写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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 18:24:52