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

Telegram bot批量上传文件后仅发送一次通知的实现位置问询

Telegram Bot 批量媒体上传统一处理实现方案

核心逻辑说明

Telegram 客户端批量上传的多份媒体资源,会被分配相同的media_group_id字段,我们可以基于该字段对同批次上传的媒体做分组收集,待整组媒体全部接收完成后再统一处理、只发送一次通知。

实现位置与步骤

你不需要修改原有的单文件处理核心逻辑,只需要在handler外层增加分组缓存、延迟判定的逻辑即可,具体实现如下:

  1. 初始化全局缓存结构,用于存储每个媒体组的消息列表和待执行的处理任务
  2. 在消息handler入口判断是否携带media_group_id,没有则走原有单文件处理逻辑
  3. 携带media_group_id的消息进入分组收集流程:将当前消息存入对应分组的缓存,每次新消息到达时重置延迟处理计时器
  4. 计时器超时后判定整组媒体接收完成,统一遍历处理所有文件,处理完成后仅发送一次通知,同时清理对应缓存

参考代码实现

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 04:45:03