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

基于aiogram的Telegram大文件批量下载失败问题及优化咨询

问题原因分析

你的机器人在处理多个大文件时出现离线、下载失败的核心原因是:

  1. 原函数中若存在同步阻塞操作(如传统FTP库的同步调用),会占用asyncio事件循环,导致机器人无法响应新请求,表现为"离线"。
  2. 大文件下载、上传属于长耗时任务,直接在消息处理器中执行会阻塞整个事件循环,无法并行处理多个任务。
  3. 多进程方案失败大概率是因为aiogram的Bot/Context实例无法跨进程直接共享,且进程间通信增加了复杂度,并非最优解。
优化方案与代码实现

以下是基于异步任务队列、异步IO优化的解决方案,完全适配aiogram的异步架构:

1. 核心优化点

  • 用异步FTP库替代同步操作,避免阻塞事件循环
  • 将长耗时任务(下载、上传)放入后台异步执行,不占用主事件循环
  • 增加并发限制,防止资源耗尽
  • 完善资源清理与错误处理

2. 完整优化代码

from aiogram import Bot, types
from aiogram.filters import MessageHandler, VIDEO, User
from aiogram.utils.markdown import hbold
import asyncio
import aioftp
from urllib.parse import urlparse
import os
import logging

# 配置日志
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

# 并发限制:同时最多处理3个任务,可根据服务器配置调整
CONCURRENCY_LIMIT = 3
semaphore = asyncio.Semaphore(CONCURRENCY_LIMIT)

# 文件名清理函数(保留你的原有实现)
def clean_filename(filename):
    return filename.replace("/", "_").replace("\\", "_").replace(":", "_") if filename else ""

# 异步FTP文件检查
async def ftp_file_check(filename, ftp_config):
    try:
        async with aioftp.Client.context(
            ftp_config["host"],
            ftp_config["user"],
            ftp_config["passwd"]
        ) as client:
            await client.stat(f"{ftp_config['remote_path']}/{filename}")
        return True
    except aioftp.PathNotFoundError:
        return False
    except Exception as e:
        logger.error(f"FTP检查失败: {str(e)}")
        return False

# 异步FTP分块上传
async def ftp_file_upload(filename, ftp_config):
    try:
        async with aioftp.Client.context(
            ftp_config["host"],
            ftp_config["user"],
            ftp_config["passwd"]
        ) as client:
            async with client.upload_stream(f"{ftp_config['remote_path']}/{filename}") as stream:
                async with open(filename, "rb") as f:
                    # 1MB分块上传,避免内存溢出
                    while chunk := await f.read(1024 * 1024):
                        await stream.write(chunk)
        return True
    except Exception as e:
        logger.error(f"FTP上传失败: {str(e)}")
        return False

# 后台任务:处理下载、FTP检查与上传
async def process_video_task(update: types.Update, bot: Bot, filename: str, file: types.File, ftp_config):
    async with semaphore:
        chat_id = update.effective_chat.id
        try:
            # 下载文件到本地
            await file.download_to_drive(filename)
            await bot.send_message(chat_id, text=f"{hbold(filename)} 下载完成", parse_mode="HTML")
            
            # FTP操作
            file_exists = await ftp_file_check(filename, ftp_config)
            if file_exists:
                await bot.send_message(chat_id, text=f"{hbold(filename)} 已存在于FTP服务器", parse_mode="HTML")
            else:
                upload_success = await ftp_file_upload(filename, ftp_config)
                if upload_success:
                    await bot.send_message(chat_id, text=f"{hbold(filename)} FTP上传成功", parse_mode="HTML")
                else:
                    await bot.send_message(chat_id, text=f"{hbold(filename)} FTP上传失败", parse_mode="HTML")
            
        except Exception as e:
            error_msg = f"{hbold(filename)} 处理出错: {str(e)}"
            await bot.send_message(chat_id, text=error_msg, parse_mode="HTML")
            logger.error(error_msg)
        finally:
            # 清理本地文件,避免磁盘占用
            if os.path.exists(filename):
                os.remove(filename)

# 消息处理器:仅负责接收请求、启动后台任务
async def downloader_video_fnc(update: types.Update, context):
    try:
        video = update.message.video
        file = await context.bot.get_file(video)
        
        # 确定文件名
        filename = clean_filename(video.file_name)
        if not filename:
            parsed_path = urlparse(file.file_path)
            filename = clean_filename(os.path.basename(parsed_path.path))
        
        # 告知用户任务启动
        await update.message.reply_text(
            f"{hbold(filename)} 任务已启动,后台处理中...",
            parse_mode="HTML",
            disable_web_page_preview=True
        )
        
        # 启动后台异步任务,不阻塞主事件循环
        ftp_config = {
            "host": "你的FTP地址",
            "user": "FTP用户名",
            "passwd": "FTP密码",
            "remote_path": "/远程存储路径"
        }
        asyncio.create_task(
            process_video_task(update, context.bot, filename, file, ftp_config)
        )
        
    except Exception as e:
        await update.message.reply_text(f"任务启动失败: {str(e)}", parse_mode="HTML")
        logger.error(f"任务启动失败: {str(e)}")

# 注册处理器
application.add_handler(MessageHandler(VIDEO, downloader_video_fnc, User(username="@XXXXXXXX")))

3. 额外优化建议

  • FTP连接池:频繁创建FTP连接会增加开销,可基于aioftp实现连接池复用连接
  • 任务状态跟踪:用字典或轻量数据库(如SQLite)记录任务ID、状态,支持用户查询任务进度
  • 文件去重:可通过文件哈希值替代文件名判断重复,避免文件名冲突导致的误判
  • 超时控制:给FTP操作、下载操作添加超时限制,防止任务无限挂起
  • 资源监控:监控磁盘空间、网络带宽,避免资源耗尽导致服务异常

内容的提问来源于stack exchange,提问作者user3130061

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 04:18:11