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

如何基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 16:05:33