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
相关产品推荐
相关产品推荐

