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

如何将Asyncio协程传入multiprocessing进程?业务场景求解

如何结合multiprocessing与asyncio实现多进程异步计算,同时保留Telegram Bot独立任务?

我开发了一款运行良好的异步计算应用,核心逻辑是无限循环运行计算,当计算满足特定条件时调用Telegram Bot的函数发送结果。但目前应用速度过慢,我打算把核心计算函数链放到多进程中运行,但始终没法让进程内的异步代码正常工作——试过multiprocessing和aioprocessing都没效果。另外必须保留Telegram Bot的独立异步任务,不能影响它的运行。

现有代码如下:

async def main():
    try:
        pairs = [[..., ...], [..., ...], [..., ...], [..., ...], [..., ...]]
        tasks_run = [func(pair) for pair in pairs]
        await asyncio.gather(*tasks_run)
    except Exception:
        logger.exception('main')


if __name__ == '__main__':
    loop = asyncio.new_event_loop()
    asyncio.set_event_loop(loop)
    tasks = [loop.create_task(dp.start_polling()), loop.create_task(main())]
    try:
        loop.run_until_complete(asyncio.gather(*tasks))
    except KeyboardInterrupt:
        pass
    finally:
        for task in tasks:
            task.cancel()
        loop.run_until_complete(loop.shutdown_asyncgens())
        loop.close()

我需要为每个pair创建独立进程,并在进程内运行协程func(pair)。想请教:

  1. 能不能结合multiprocessing与asyncio实现需求?如果可以该怎么操作?
  2. 如果不行,有什么替代方案?
  3. 主入口同时运行Telegram Bot轮询和计算任务,添加多进程会不会破坏二者的交互?

可以结合multiprocessing与asyncio实现,核心是进程隔离+跨进程通信

1. 核心思路:每个子进程独立运行asyncio事件循环

asyncio的事件循环是进程内资源,子进程无法复用主进程的循环。所以要把异步计算逻辑包装成能在子进程内启动循环的函数,再用multiprocessing.Process启动每个进程。

2. 关键:跨进程传递计算结果给主进程的Telegram Bot

子进程不能直接调用主进程的Bot对象(跨进程无法共享异步状态),必须用**进程间通信(IPC)**传递消息。推荐用multiprocessing.Queue:子进程满足条件时把结果放进队列,主进程专门开一个异步任务监听队列,收到消息后调用Bot发送。

3. 具体实现代码示例

第一步:封装子进程的异步运行逻辑

import multiprocessing
import asyncio

# 子进程的包装函数,负责启动本地事件循环并运行计算协程
def process_worker(pair, result_queue):
    # 每个子进程创建自己的事件循环
    asyncio.run(run_func_with_queue(pair, result_queue))

async def run_func_with_queue(pair, result_queue):
    # 执行原有的计算逻辑
    result = await func(pair)
    if 满足特定条件:
        # 将可序列化的结果放入队列(不能传递异步对象)
        result_queue.put(result)

第二步:修改主函数,启动多进程+监听队列

# 保存全局进程列表,用于中断时清理
processes = []

async def main(result_queue):
    global processes
    try:
        pairs = [[..., ...], [..., ...], [..., ...], [..., ...], [..., ...]]
        # 为每个pair创建子进程
        for pair in pairs:
            p = multiprocessing.Process(target=process_worker, args=(pair, result_queue))
            processes.append(p)
            p.start()
        
        # 无限循环计算的场景,不需要join,让进程持续运行
        while True:
            await asyncio.sleep(3600)
    except Exception:
        logger.exception('main')

# 主进程监听队列的异步任务,负责调用Telegram Bot发送结果
async def listen_queue(result_queue, dp):
    while True:
        # 用asyncio.to_thread避免阻塞主事件循环
        result = await asyncio.to_thread(result_queue.get)
        # 调用你的Telegram Bot发送逻辑
        await dp.bot.send_message(chat_id=你的聊天ID, text=str(result))

第三步:修改主入口,整合所有任务

if __name__ == '__main__':
    # 创建跨进程队列
    result_queue = multiprocessing.Queue()
    
    loop = asyncio.new_event_loop()
    asyncio.set_event_loop(loop)
    
    # 三个核心任务:Bot轮询、多进程计算、队列监听
    tasks = [
        loop.create_task(dp.start_polling()),
        loop.create_task(main(result_queue)),
        loop.create_task(listen_queue(result_queue, dp))
    ]
    
    try:
        loop.run_until_complete(asyncio.gather(*tasks))
    except KeyboardInterrupt:
        # 终止所有子进程
        for p in processes:
            p.terminate()
            p.join()
    finally:
        for task in tasks:
            task.cancel()
        loop.run_until_complete(loop.shutdown_asyncgens())
        loop.close()

4. 替代方案:用concurrent.futures.ProcessPoolExecutor

如果觉得手动管理进程麻烦,可以用ProcessPoolExecutor自动管理进程池。注意Executor中的函数需为同步函数,或在函数内部启动asyncio循环:

from concurrent.futures import ProcessPoolExecutor

async def main(result_queue, dp):
    pairs = [[..., ...], [..., ...], [..., ...], [..., ...], [..., ...]]
    
    # 启动进程池,数量等于pair的数量
    with ProcessPoolExecutor(max_workers=len(pairs)) as executor:
        # 提交所有计算任务
        futures = [executor.submit(process_worker, pair, result_queue) for pair in pairs]
        # 同时启动队列监听任务,等待所有计算任务完成
        await asyncio.gather(
            listen_queue(result_queue, dp),
            *[asyncio.wrap_future(f) for f in futures]
        )

关于Telegram Bot交互的问题

不会破坏交互。只要遵循以下规则:

  • Telegram Bot的所有操作都在主进程的事件循环中执行,子进程只负责计算,不直接操作Bot对象。
  • 子进程通过队列传递结果,主进程的监听任务统一处理Bot发送逻辑。
    这样Bot的轮询任务和计算进程完全隔离,互不干扰。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 18:27:52