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
相关产品推荐
相关产品推荐

