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

如何将Python Aiogram与多进程结合?解决调度器阻塞问题

问题描述

重写解析器用于Telegram Bot,需要从函数中启动多进程执行任务,但multiprocessing进程运行时,Aiogram的Dispatcher无法响应新聊天消息,调度器处于闲置状态。

问题代码
from aiogram import Dispatcher, Bot, executor
from multiprocessing import Process
from time import sleep

bot_token = "token"
bot = Bot(token=bot_token)
dp = Dispatcher(bot)
allowed_users = [My_user_id]


class work:
    @staticmethod
    def reparse(user_id, data):
        sleep(10)  # 模拟任务耗时
        return

    @staticmethod
    async def parsing(user_id, message):
        await message.answer('Parsing started')
        some_data = [124215, 123543, 346457, 347347]
        processes = []
        for i in range(4):
            args = [user_id, some_data[i]]
            _process = Process(target=work.reparse, args=(*args,))
            processes.append(_process)
            _process.start()
        for process in processes:
            process.join()
        await message.answer('Work completed')
        return


if __name__ == '__main__':
    @dp.message_handler(commands=['check_dispatcher'])
    async def check_dispatcher(message):
        await message.answer('ANSWER')


    @dp.message_handler(commands=['work_start'])
    async def parser(message):
        await work.parsing(message.chat.id, message)


    print('STARTED')
    executor.start_polling(dp, skip_updates=True)
已尝试方案(均无效)
  • 使用process.join(),调度器在进程运行时仍闲置;去掉join()改用process.is_alive()轮询状态,仍无效。
  • 尝试通过asyncio启动进程,未解决阻塞问题。
  • 尝试threading.Thread结合asyncio.run,用法错误,未解决问题。
问题根源

核心问题是同步阻塞调用卡住了asyncio事件循环:Process.join()是同步阻塞方法,会让aiogram依赖的asyncio事件循环无法处理新的消息更新,直到所有进程完成。哪怕用循环轮询is_alive(),也会持续占用事件循环,导致Dispatcher无法响应新消息。

解决方案

使用asyncio.to_thread将同步的join操作放到后台线程执行,让asyncio事件循环可以继续处理其他任务(比如新的聊天消息)。修改后的parsing方法如下:

import asyncio  # 需要导入asyncio

class work:
    @staticmethod
    def reparse(user_id, data):
        sleep(10)  # 模拟任务耗时
        return

    @staticmethod
    async def parsing(user_id, message):
        await message.answer('Parsing started')
        some_data = [124215, 123543, 346457, 347347]
        processes = []
        for i in range(4):
            args = [user_id, some_data[i]]
            _process = Process(target=work.reparse, args=(*args,))
            processes.append(_process)
            _process.start()
        
        # 将每个进程的join操作放到后台线程,异步等待所有完成
        await asyncio.gather(*[asyncio.to_thread(process.join) for process in processes])
        
        await message.answer('Work completed')
        return
方案说明
  • asyncio.to_thread(func)会把同步函数func放到单独的线程中执行,asyncio事件循环不会被阻塞,可以继续处理Dispatcher的消息更新。
  • asyncio.gather(*tasks)用于异步等待所有后台线程的join操作完成,确保所有进程结束后再发送完成通知。
补充注意事项
  • 确保导入asyncio模块,否则无法使用to_thread和gather。
  • 如果需要在进程中与Bot交互,不要直接在子进程中使用aiogram的Bot实例,多进程环境下Bot实例不能共享,建议通过队列传递结果,再由主进程发送消息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 22:47:03