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

Python中AsyncIO与Process Pool结合使用的实现方案问询

结合AsyncIO与多进程的实现方案

你的需求是完全可实现的,这也是Python异步生态里利用多核CPU的标准实践之一,整体逻辑和你描述的完全吻合。

核心实现思路

  • 跨进程队列使用multiprocessing.Queue(或者JoinableQueue)做任务分发和结果回收,这类队列是进程安全的,支持多进程同时读写
  • 每个工作进程启动时初始化独立的AsyncIO事件循环,进程内部所有任务都由该事件循环调度并发执行
  • 主线程的AsyncIO事件循环仅负责接收外部请求、推送任务到共享队列、拉取结果回传给请求方,不处理具体业务逻辑

最简实现示例

import asyncio
import multiprocessing
from typing import Callable, Any, Coroutine

# 工作进程入口
def worker_process(task_queue: multiprocessing.Queue, result_queue: multiprocessing.Queue):
    # 每个进程初始化自己的事件循环
    loop = asyncio.new_event_loop()
    asyncio.set_event_loop(loop)

    async def _worker_routine():
        while True:
            # 从任务队列取任务,阻塞调用放到线程池执行避免卡住事件循环
            task = await loop.run_in_executor(None, task_queue.get)
            if task is None: # 收到终止信号退出
                task_queue.task_done()
                break
            # 任务格式: (任务ID, 异步处理函数, 位置参数, 关键字参数)
            task_id, func, args, kwargs = task
            try:
                result = await func(*args, **kwargs)
                result_queue.put((task_id, True, result))
            except Exception as e:
                result_queue.put((task_id, False, str(e)))
            finally:
                task_queue.task_done()
    
    loop.run_until_complete(_worker_routine())
    loop.close()

# 主线程的任务调度器
class AsyncProcessPool:
    def __init__(self, worker_num: int = multiprocessing.cpu_count()):
        self.task_queue = multiprocessing.JoinableQueue()
        self.result_queue = multiprocessing.Queue()
        self.workers = []
        # 启动工作进程
        for _ in range(worker_num):
            p = multiprocessing.Process(
                target=worker_process,
                args=(self.task_queue, self.result_queue)
            )
            p.daemon = True
            p.start()
            self.workers.append(p)
        # 启动结果拉取协程
        self._pending_tasks = {}
        self._next_task_id = 0
        asyncio.create_task(self._pull_result_routine())
    
    async def _pull_result_routine(self):
        loop = asyncio.get_running_loop()
        while True:
            task_id, success, result = await loop.run_in_executor(None, self.result_queue.get)
            if task_id in self._pending_tasks:
                fut = self._pending_tasks.pop(task_id)
                if success:
                    fut.set_result(result)
                else:
                    fut.set_exception(RuntimeError(result))
    
    async def submit(self, func: Callable[..., Coroutine], *args, **kwargs) -> Any:
        fut = asyncio.get_running_loop().create_future()
        task_id = self._next_task_id
        self._next_task_id += 1
        self._pending_tasks[task_id] = fut
        self.task_queue.put((task_id, func, args, kwargs))
        return await fut
    
    async def shutdown(self):
        # 给所有工作进程发终止信号
        for _ in range(len(self.workers)):
            self.task_queue.put(None)
        await self.task_queue.join()
        for p in self.workers:
            p.join()

# 测试用的异步任务
async def demo_task(x: int) -> int:
    await asyncio.sleep(1)
    return x * x

# 测试入口
async def main():
    pool = AsyncProcessPool(worker_num=4)
    # 同时提交10个任务,4个进程每个进程并发处理,总耗时约3秒(10/4向上取整)
    tasks = [pool.submit(demo_task, i) for i in range(10)]
    results = await asyncio.gather(*tasks)
    print(results)
    await pool.shutdown()

if __name__ == "__main__":
    asyncio.run(main())

注意事项

  • 跨进程传递的函数、参数、结果必须支持pickle序列化,否则会出现入队失败的问题
  • 如果业务中有大量小任务,可以给每个工作进程设置单进程并发上限,避免事件循环调度压力过大导致延迟升高
  • 可以根据业务需要添加任务超时、重试、队列长度限制等逻辑,提升稳定性

内容的提问来源于stack exchange,提问作者user688661

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 00:24:01