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

Telethon多进程批量下载大文件失效,寻求优化解决方案

Telegram大文件批量下载提速与多进程并发问题解决

问题背景

从Telegram群组批量下载1GB以上文件,单线程使用Telethon下载速度仅5-6Mbps,已安装cryptg但无明显改善;尝试用multiprocessing实现多文件并发下载时无法正常运行,附示例代码寻求解决思路。

原错误代码

from telethon import TelegramClient
from multiprocessing import Pool
from datetime import datetime
import time
import threading
import os


# Remember to use your own values from my.telegram.org!
api_id = HiddenByME
api_hash = 'HiddenByME'
bot_token= 'HiddenByME'

client = TelegramClient('anon', api_id, api_hash)
bot= TelegramClient('bot', api_id, api_hash).start(bot_token=bot_token)

movieschanneltofix = HiddenByME
movieschannelfixed = HiddenByME

downloadpath = "media/"

# Printing download progress
def callbackdownload(current, total):
    global start_time
    global end_time
    delta = (datetime.now()- start_time).total_seconds()
    bytes_remaining = total - current
    speed = round(((current / 1024) / 1024) / delta, 2)
    seconds_left = round(((bytes_remaining / 1024) / 1024) / float(speed), 2)
    print('Downloaded', round(current/1000000,2), 'out of', round(total/1000000,2),'Megabytes {:.2%}'.format(current / total), "->", "Trascorso:", round(delta/60, 2), " Minuti * Rimanente:", int(seconds_left), "Secondi * Velocità:", round(speed, 2), "Mbps",  end='\r')
# Printing Upload Progress
def callbackupload(current, total):
    global start_time
    global end_time
    delta = (datetime.now()- start_time).total_seconds()
    bytes_remaining = total - current
    speed = round(((current / 1024) / 1024) / delta, 2)
    seconds_left = round(((bytes_remaining / 1024) / 1024) / float(speed), 2)
    print('Uploaded', round(current/1000000,2), 'out of', round(total/1000000,2),'Megabytes {:.2%}'.format(current / total), "->", "Trascorso:", round(delta/60, 2), " Minuti * Rimanente:", int(seconds_left), "Secondi * Velocità:", round(speed, 2), "Mbps",  end='\r')


async def moviefix(results):
    global end_time
    global start_time
    print("Fixing Message ID: ", results.split(';')[0], "\n", "OldName: ", results.split(';')[1], "\n", "NewName: ", results.split(';')[2])

    # Get Message and Download Media Renaming It
    downpath = downloadpath+results.split(';')[2]
    start_time = datetime.now()
    print("Download Avviato:", start_time.strftime("%d/%m/%Y, %H:%M:%S"))
    path = await client.download_media(await client.get_messages(movieschanneltofix, ids=int(results.split(';')[0])), downpath, progress_callback=callbackdownload)
    print('Download Completato!!! ', path)
    time.sleep(1)

    #Upload Renamed File
    start_time = datetime.now()
    print("Upload Avviato:", start_time.strftime("%d/%m/%Y, %H:%M:%S"))
    result = await bot.send_file(movieschannelfixed, downpath, force_document=True, progress_callback=callbackupload)
    print('Upload Completato!!! ', result.file.name)
    time.sleep(1)

    #Delete Uploaded File
    os.remove(downpath)
    time.sleep(1)

    #Save Fixed File in DB
    with open('data/lastmoviefixed.txt', 'a') as f:
        f.write(results+"\n")
    time.sleep(1)
    time.sleep(20)
    return results


#    print (results, "PID:", os.getpid())
#    with open('data/lastmoviefixed.txt', 'a') as f:
#        f.write(results+"\n")
#        time.sleep(1)
#    time.sleep(10)
#    return results

if __name__ == "__main__":
    global listtodo
    listtodo = set(open('data/listtofix.txt', 'r', encoding="utf8").read().splitlines()) - set(open('data/lastmoviefixed.txt', encoding="utf8").read().splitlines())
    with Pool(1) as p:
        result = p.map_async(client.loop.run_until_complete(moviefix), sorted(listtodo))
        while not result.ready():
            time.sleep(1)
        result = result.get()
        p.terminate()
        p.join()
    print("Done", results)

核心问题分析

原代码的多进程实现完全错误:

  • Telethon客户端是协程对象,不能跨进程共享,每个进程必须独立初始化客户端
  • map_async的第一个参数应为可调用函数,而非直接执行client.loop.run_until_complete(moviefix)
  • 全局start_time在多任务场景下会被覆盖,导致进度计算混乱

解决方案

1. 改用多协程替代多进程

Telethon基于asyncio开发,多协程更适合IO密集型的下载/上传任务,无需引入multiprocessing,避免进程间通信和客户端共享问题。

2. 重构后的完整代码

from telethon import TelegramClient
from datetime import datetime
import asyncio
import os

# 替换为你的实际参数
api_id = HiddenByME
api_hash = 'HiddenByME'
bot_token= 'HiddenByME'

movieschanneltofix = HiddenByME
movieschannelfixed = HiddenByME
downloadpath = "media/"

# 封装进度回调(避免全局变量冲突)
def create_progress_callback(task_name):
    start_time = datetime.now()
    def callback(current, total):
        delta = (datetime.now() - start_time).total_seconds()
        if delta == 0:
            return
        # 计算进度与速度
        speed_mb = round((current / 1024 / 1024) / delta, 2)
        current_mb = round(current / 1000000, 2)
        total_mb = round(total / 1000000, 2)
        progress = current / total * 100
        print(f'{task_name}: {current_mb}/{total_mb}MB ({progress:.2f}%) | 速度: {speed_mb}Mbps | 已耗时: {round(delta/60,2)}分钟', end='\r')
    return callback

async def moviefix(client, bot, task_item):
    msg_id, old_name, new_name = task_item.split(';')
    downpath = os.path.join(downloadpath, new_name)
    
    print(f"\n处理消息ID: {msg_id}\n原文件名: {old_name}\n新文件名: {new_name}")

    # 下载文件
    print(f"开始下载: {datetime.now().strftime('%d/%m/%Y, %H:%M:%S')}")
    path = await client.download_media(
        await client.get_messages(movieschanneltofix, ids=int(msg_id)),
        downpath,
        progress_callback=create_progress_callback(f"下载 {new_name}")
    )
    print(f'\n下载完成: {path}')

    # 上传文件
    print(f"开始上传: {datetime.now().strftime('%d/%m/%Y, %H:%M:%S')}")
    upload_result = await bot.send_file(
        movieschannelfixed,
        downpath,
        force_document=True,
        progress_callback=create_progress_callback(f"上传 {new_name}")
    )
    print(f'\n上传完成: {upload_result.file.name}')

    # 清理本地文件
    os.remove(downpath)

    # 记录已处理项
    with open('data/lastmoviefixed.txt', 'a', encoding='utf8') as f:
        f.write(task_item + "\n")

    await asyncio.sleep(2)  # 避免请求过于频繁触发限流
    return task_item

async def main():
    # 独立初始化客户端(每个协程共用同一客户端即可,无需多进程)
    client = TelegramClient('anon', api_id, api_hash)
    bot = TelegramClient('bot', api_id, api_hash)
    await client.start()
    await bot.start(bot_token=bot_token)

    # 获取待处理任务列表
    with open('data/listtofix.txt', 'r', encoding='utf8') as f:
        all_tasks = set(f.read().splitlines())
    with open('data/lastmoviefixed.txt', 'r', encoding='utf8') as f:
        done_tasks = set(f.read().splitlines())
    todo_tasks = sorted(all_tasks - done_tasks)

    # 控制并发数(根据网络带宽调整,建议2-5)
    concurrency_limit = 3
    semaphore = asyncio.Semaphore(concurrency_limit)

    # 绑定并发限制的任务包装函数
    async def bounded_task(task):
        async with semaphore:
            return await moviefix(client, bot, task)

    # 批量执行任务
    tasks = [bounded_task(task) for task in todo_tasks]
    await asyncio.gather(*tasks)

    # 关闭客户端
    await client.disconnect()
    await bot.disconnect()
    print("\n所有任务处理完成")

if __name__ == "__main__":
    asyncio.run(main())

3. 额外提速优化建议

  • 验证cryptg生效:在代码开头加入print(telethon.crypto.cryptg is not None),返回True则说明已启用加速
  • 调整并发数:根据你的网络带宽和Telegram服务器限制,合理设置concurrency_limit,过高会触发限流
  • 添加错误处理:给moviefix函数添加try-except块,处理下载/上传失败的情况,避免单个任务崩溃导致整体终止
  • 优化存储IO:确保downloadpath所在磁盘为SSD,减少文件读写耗时

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 17:05:03