如何基于python-telegram-bot Webhooks同时处理多用户消息?
aiogram Webhook 实现多消息并行处理方案
可行性说明
Webhook 模式完全支持并行处理多用户消息,核心是利用 aiogram 基于**异步IO(asyncio)**的特性,配合异步任务调度或线程/进程池,将耗时逻辑从主请求处理流程中剥离,避免串行阻塞。
具体实现方案
方案1:异步任务调度(适配IO密集型场景)
适合你的10秒等待类IO密集型任务,通过asyncio.create_task将耗时逻辑转为后台异步任务,主请求处理线程立即释放,可同时接收并处理新的用户消息。
示例代码:
import asyncio from aiogram import Bot, Dispatcher, types from aiogram.types import Update from aiohttp import web API_TOKEN = '你的机器人Token' bot = Bot(token=API_TOKEN) dp = Dispatcher(bot) # 封装耗时的异步处理逻辑 async def process_message_task(message: types.Message): await asyncio.sleep(10) # 模拟10秒耗时操作 await message.reply("处理完成!") @dp.message_handler() async def handle_message(message: types.Message): # 立即回复用户,避免等待 await message.reply("正在处理,请稍候...") # 提交异步任务,不阻塞当前请求 asyncio.create_task(process_message_task(message)) async def handle_webhook(request: web.Request): update = Update(**await request.json()) await dp.process_update(update) return web.Response() async def main(): app = web.Application() app.add_routes([web.post('/webhook', handle_webhook)]) # 配置Webhook地址(需确保域名/端口可被Telegram访问) await bot.set_webhook(url='https://你的域名/webhook') runner = web.AppRunner(app) await runner.setup() site = web.TCPSite(runner, '0.0.0.0', 8080) await site.start() await asyncio.Event().wait() if __name__ == '__main__': asyncio.run(main())
此方案中,多个用户的10秒任务会在asyncio事件循环中并行执行(IO等待时自动切换任务),总耗时不会随用户数量叠加。
方案2:线程/进程池(适配CPU密集型场景)
如果你的耗时逻辑是CPU密集型(而非单纯IO等待),可通过concurrent.futures创建线程池或进程池,实现真正的并行执行,和你第一个示例的worker并行逻辑完全匹配。
示例代码(线程池版):
import asyncio import time from concurrent.futures import ThreadPoolExecutor from aiogram import Bot, Dispatcher, types from aiogram.types import Update from aiohttp import web API_TOKEN = '你的机器人Token' bot = Bot(token=API_TOKEN) dp = Dispatcher(bot) # 创建线程池,设置最大并行数(对应你第一个示例的4个worker) executor = ThreadPoolExecutor(max_workers=4) # 耗时的同步处理函数 def heavy_processing(chat_id: int, message_id: int): time.sleep(10) # 模拟10秒CPU/IO任务 return chat_id, message_id async def process_with_pool(message: types.Message): # 在线程池中执行耗时任务 chat_id, msg_id = await asyncio.get_event_loop().run_in_executor( executor, heavy_processing, message.chat.id, message.message_id ) await bot.send_message(chat_id=chat_id, text=f"消息 {msg_id} 处理完成!") @dp.message_handler() async def handle_message(message: types.Message): await message.reply("正在处理,请稍候...") # 提交线程池任务,异步执行 asyncio.create_task(process_with_pool(message)) async def handle_webhook(request: web.Request): update = Update(**await request.json()) await dp.process_update(update) return web.Response() async def main(): app = web.Application() app.add_routes([web.post('/webhook', handle_webhook)]) await bot.set_webhook(url='https://你的域名/webhook') runner = web.AppRunner(app) await runner.setup() site = web.TCPSite(runner, '0.0.0.0', 8080) await site.start() await asyncio.Event().wait() if __name__ == '__main__': try: asyncio.run(main()) finally: executor.shutdown()
若使用进程池,需注意Bot实例无法跨进程共享,需在子进程中重新初始化,或仅传递chat_id、message_id等必要参数。
关键注意事项
- 必须使用异步Web框架(aiogram默认搭配aiohttp,满足要求),同步Web框架会导致请求串行处理。
- 生产环境建议配合Nginx等反向代理,确保Webhook请求稳定转发。
- 避免在异步任务中阻塞事件循环,CPU密集型任务优先用线程/进程池。
内容的提问来源于stack exchange,提问作者Tomas Angelo
相关产品推荐
相关产品推荐

