如何使用Pyrogram并行处理Telegram消息
Pyrogram实现多消息并行处理的方法
Pyrogram默认会按消息接收顺序逐个处理,要实现并行处理核心是把耗时的业务逻辑从消息处理的同步流程中剥离,交给异步事件池并行执行。以下是两种实用方案:
方案一:直接用asyncio.create_task()并行执行任务
将耗时的异步函数包装成独立任务,丢给事件循环后台执行,这样当前消息处理流程不会被阻塞,能立即处理下一条消息。
修改后的代码示例:
from pyrogram import Client import asyncio app = Client(client_name, api_id=api_id, api_hash=api_hash) async def some_other_function(): # 模拟耗时操作(如IO请求、复杂计算) await asyncio.sleep(5) print("耗时任务执行完成") async def handle_messages(client, message): # 把耗时任务异步提交到后台 asyncio.create_task(some_other_function()) # 可立即回复用户,无需等待耗时任务结束 await message.reply("消息已接收,正在处理中") app.run()
方案二:用信号量限制并发数
如果短时间内收到大量消息,无限制创建任务可能导致资源耗尽。可以用asyncio.Semaphore控制同时运行的最大任务数:
from pyrogram import Client import asyncio # 设置最大并发任务数为5 max_concurrent_tasks = asyncio.Semaphore(5) app = Client(client_name, api_id=api_id, api_hash=api_hash) async def some_other_function(): async with max_concurrent_tasks: # 模拟耗时操作 await asyncio.sleep(5) print("耗时任务执行完成") async def handle_messages(client, message): asyncio.create_task(some_other_function()) await message.reply("消息已接收,正在处理中") app.run()
注意事项
- 确保
some_other_function是异步函数,如果是同步耗时操作,需要用loop.run_in_executor把它放到线程池执行,避免阻塞事件循环。 - 如果需要等待任务完成后再回复用户,可以把信号量逻辑放到
handle_messages里,但这样会阻塞当前回复流程,仅适合必须等待处理结果的场景。
内容的提问来源于stack exchange,提问作者understack
相关产品推荐
相关产品推荐

