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

