Telegram bot批量上传文件后仅发送一次通知的实现位置问询
Telegram Bot 批量媒体上传统一处理实现方案
核心逻辑说明
Telegram 客户端批量上传的多份媒体资源,会被分配相同的media_group_id字段,我们可以基于该字段对同批次上传的媒体做分组收集,待整组媒体全部接收完成后再统一处理、只发送一次通知。
实现位置与步骤
你不需要修改原有的单文件处理核心逻辑,只需要在handler外层增加分组缓存、延迟判定的逻辑即可,具体实现如下:
- 初始化全局缓存结构,用于存储每个媒体组的消息列表和待执行的处理任务
- 在消息handler入口判断是否携带
media_group_id,没有则走原有单文件处理逻辑 - 携带
media_group_id的消息进入分组收集流程:将当前消息存入对应分组的缓存,每次新消息到达时重置延迟处理计时器 - 计时器超时后判定整组媒体接收完成,统一遍历处理所有文件,处理完成后仅发送一次通知,同时清理对应缓存
参考代码实现
import asyncio from typing import Dict, List, Tuple # 全局缓存:key为媒体组ID,value为(该组所有消息列表, 延迟处理任务对象) media_group_cache: Dict[str, Tuple[List[types.Message], asyncio.Task | None]] = {} async def process_media_group(media_group_id: str, user_id: int, username: str): """统一处理已收集完成的整组媒体""" messages, _ = media_group_cache.pop(media_group_id) try: # 批量下载所有文件 for msg in messages: files = await bot.get_file(msg.photo[-1].file_id) save_path = f'data/{user_id}/{files.file_path[-12:]}' await bot.download_file_by_id(msg.photo[-1].file_id, save_path) logger.info('{} {}', username, save_path) # 全部处理完成后仅发送一次通知 await messages[0].answer('Все сохранено', reply_markup=inkb) except Exception as e: await messages[0].answer(f'处理失败: {e}') async def handle_docs_photo(message: types.Message): user_id = message.from_user.id username = message.from_user.username media_group_id = message.media_group_id # 无媒体组ID,为单张上传,走原有逻辑 if not media_group_id: try: files = await bot.get_file(message.photo[-1].file_id) save_path = f'data/{user_id}/{files.file_path[-12:]}' await bot.download_file_by_id(message.photo[-1].file_id, save_path) logger.info('{} {}', username, save_path) await message.answer('Все сохранено', reply_markup=inkb) except Exception as e: await message.answer(f'处理失败: {e}') return # 批量上传,走分组收集逻辑 if media_group_id not in media_group_cache: media_group_cache[media_group_id] = ([message], None) else: media_group_cache[media_group_id][0].append(message) # 取消之前的延迟任务,重置计时 old_task = media_group_cache[media_group_id][1] if old_task: old_task.cancel() # 新建延迟任务,1秒内无同组新消息则触发批量处理,延迟时长可根据实际情况调整 new_task = asyncio.create_task(process_media_group(media_group_id, user_id, username)) media_group_cache[media_group_id] = (media_group_cache[media_group_id][0], new_task)
注意事项
- 延迟处理的时长可根据网络情况调整,一般0.5~2秒即可覆盖同组媒体的到达间隔
- 多实例部署场景下,可将内存缓存替换为Redis等共享存储,保证跨实例的同组媒体可被正确收集
- 可额外增加定时清理逻辑,避免异常场景下的缓存残留占用内存
内容的提问来源于stack exchange,提问作者lonelywolf
相关产品推荐
相关产品推荐

