如何优化Python中asyncio+Telethon同步循环性能并解决限流问题
Telegram群组用户数据导出优化问题
我用Python的Telethon库从Telegram群组获取用户姓名、用户名及简介并导出到JSON文件。运行时发现,约30次请求后获取操作速度骤降,定位到是遍历用户同步调用GetFullUserRequest导致的性能瓶颈。
尝试用asyncio并发请求,80人群组正常,但用户量更大时出现限流错误:
Server response had invalid buffer: Invalid response buffer (HTTP code 429)Server indicated flood error at transport level: Invalid response buffer (HTTP code 429)
随后用asyncio.Semaphore(10)限制并发数,问题仍未解决,且循环结束后还有长时间延迟。请问问题原因是什么,如何优化代码提升效率?
运行环境:Python 3.11.6
原始代码
import configparser import json import asyncio from telethon.tl.functions.users import GetFullUserRequest from telethon import TelegramClient from telethon.errors import SessionPasswordNeededError from telethon.tl.functions.channels import GetParticipantsRequest from telethon.tl.types import ChannelParticipantsSearch from telethon.tl.types import ( PeerChannel ) # Reading Configs config = configparser.ConfigParser() config.read("config.ini") # Setting configuration values api_id = config['Telegram']['api_id'] api_hash = config['Telegram']['api_hash'] api_hash = str(api_hash) phone = config['Telegram']['phone'] username = config['Telegram']['username'] # Create the client and connect client = TelegramClient(username, api_id, api_hash) async def main(phone): await client.start() print("Client Created") # Ensure you're authorized if await client.is_user_authorized() == False: await client.send_code_request(phone) try: await client.sign_in(phone, input('Enter the code: ')) except SessionPasswordNeededError: await client.sign_in(password=input('Password: ')) me = await client.get_me() user_input_channel = input("enter entity(telegram URL or entity id):") if user_input_channel.isdigit(): entity = PeerChannel(int(user_input_channel)) else: entity = user_input_channel my_channel = await client.get_entity(entity) offset = 0 limit = 100 all_participants = [] count = 0 while True: participants = await client(GetParticipantsRequest(my_channel, ChannelParticipantsSearch(''), offset, limit,hash=0)) if not participants.users or offset >= 100: break all_participants.extend(participants.users) count +=1 print(len(participants.users)) offset += len(participants.users) print("finished") all_user_details = [] for participant in all_participants: user_full = await client(GetFullUserRequest(participant.id)) all_user_details.append({ "id": user_full.full_user.id, "bio": user_full.full_user.about }) with open('user_data.json', 'w') as outfile: json.dump(all_user_details, outfile) with client: client.loop.run_until_complete(main(phone))
优化尝试后的代码片段
# ... (other imports and code) api_semaphore = asyncio.Semaphore(10) #Updated line async def main(phone): await client.start() print("Client Created") # Ensure you're authorized if await client.is_user_authorized() == False: await client.send_code_request(phone) try: await client.sign_in(phone, input('Enter the code: ')) except SessionPasswordNeededError: await client.sign_in(password=input('Password: ')) me = await client.get_me() user_input_channel = input("enter entity(telegram URL or entity id):") if user_input_channel.isdigit(): entity = PeerChannel(int(user_input_channel)) else: entity = user_input_channel my_channel = await client.get_entity(entity) # Fetch all participants offset = 0 limit = 100 all_participants = [] count = 0 while True: participants = await client(GetParticipantsRequest(my_channel, ChannelParticipantsSearch(''), offset, limit,hash=0)) if not participants.users or offset >= 1000: break all_participants.extend(participants.users) count +=1 print(len(participants.users)) offset += len(participants.users) print("finished") # Process the participants as needed all_user_details = [] tasks = [] for participant in all_participants: async with api_semaphore: # This will prepare all requests and let them be ready to go tasks.append(asyncio.create_task(client(GetFullUserRequest(participant.id)))) #Updated line # ... (rest of the code)
问题原因分析
- Telegram API严格限流:Telegram对
GetFullUserRequest这类用户详情接口的请求频率限制极严,即便用Semaphore控制并发数,短时间内密集发起请求仍会触发429限流。 - Semaphore使用逻辑错误:你当前的写法是在循环中逐个创建任务时加锁,导致所有任务几乎同时进入等待队列,本质还是短时间内发起大量请求,没起到真正的速率控制作用。
- 未处理内置重试延迟:Telethon本身会对限流请求自动重试,但重试等待会累积,导致循环结束后仍有长时间的后台处理延迟。
优化方案
1. 正确控制请求速率(Semaphore+固定延迟)
不要只靠Semaphore限制并发,给每个请求添加固定间隔,避免短时间内请求密集:
async def safe_get_full_user(user_id, semaphore, min_delay=0.3): async with semaphore: await asyncio.sleep(min_delay) return await client(GetFullUserRequest(user_id))
2. 批量创建任务并统一等待
修改任务创建逻辑,用asyncio.gather批量处理,避免循环内逐个触发:
# 替换原任务处理部分 all_user_details = [] semaphore = asyncio.Semaphore(5) # 降低并发数至5,更贴合Telegram限流规则 tasks = [ safe_get_full_user(participant.id, semaphore) for participant in all_participants ] # 等待所有任务完成,同时捕获异常 results = await asyncio.gather(*tasks, return_exceptions=True) # 过滤异常并整理数据 for result in results: if isinstance(result, Exception): print(f"获取用户信息失败: {result}") continue user = result.full_user all_user_details.append({ "id": user.id, "first_name": user.first_name or "", "last_name": user.last_name or "", "username": user.username or "", "bio": user.about or "" })
3. 利用Telethon批量接口减少请求次数
Telethon的client.get_users支持批量获取用户基本信息,可大幅减少请求量。如果仅需姓名、用户名,直接用这个接口即可;若必须获取简介,再按需调用GetFullUserRequest:
# 批量获取用户基本信息 user_ids = [p.id for p in all_participants] users = await client.get_users(user_ids) # 按需补充简介信息(仅对需要的用户调用GetFullUserRequest)
4. 修复原代码其他问题
- 移除
offset >= 100这类不合理的终止条件,仅保留not participants.users作为停止判断。 - 导出JSON时添加
ensure_ascii=False和indent=2,保证输出中文可读且格式美观。
完整优化后代码
import configparser import json import asyncio from telethon.tl.functions.users import GetFullUserRequest from telethon import TelegramClient from telethon.errors import SessionPasswordNeededError, FloodWaitError from telethon.tl.functions.channels import GetParticipantsRequest from telethon.tl.types import ChannelParticipantsSearch, PeerChannel # 读取配置 config = configparser.ConfigParser() config.read("config.ini") api_id = config['Telegram']['api_id'] api_hash = str(config['Telegram']['api_hash']) phone = config['Telegram']['phone'] username = config['Telegram']['username'] client = TelegramClient(username, api_id, api_hash) async def safe_get_full_user(user_id, semaphore, min_delay=0.3): """带速率控制和异常重试的用户详情获取函数""" async with semaphore: await asyncio.sleep(min_delay) try: return await client(GetFullUserRequest(user_id)) except FloodWaitError as e: # 处理Telegram明确的限流等待 print(f"触发限流,等待{e.seconds}秒") await asyncio.sleep(e.seconds) return await client(GetFullUserRequest(user_id)) except Exception as e: print(f"获取用户{user_id}信息失败: {str(e)}") return None async def main(phone): await client.start() print("客户端已创建") # 授权验证 if not await client.is_user_authorized(): await client.send_code_request(phone) try: await client.sign_in(phone, input('输入验证码: ')) except SessionPasswordNeededError: await client.sign_in(password=input('输入密码: ')) # 获取群组实体 user_input_channel = input("输入群组链接或ID: ") if user_input_channel.isdigit(): entity = PeerChannel(int(user_input_channel)) else: entity = user_input_channel my_channel = await client.get_entity(entity) # 批量获取群组所有参与者 offset = 0 limit = 100 all_participants = [] while True: participants = await client( GetParticipantsRequest( my_channel, ChannelParticipantsSearch(''), offset, limit, hash=0 ) ) if not participants.users: break all_participants.extend(participants.users) offset += len(participants.users) print(f"已获取{len(all_participants)}个用户") # 并发获取用户详情 semaphore = asyncio.Semaphore(5) tasks = [ safe_get_full_user(p.id, semaphore) for p in all_participants ] results = await asyncio.gather(*tasks) # 整理数据 all_user_details = [] for result in results: if not result: continue user = result.full_user all_user_details.append({ "id": user.id, "first_name": user.first_name or "", "last_name": user.last_name or "", "username": user.username or "", "bio": user.about or "" }) # 导出JSON with open('user_data.json', 'w', encoding='utf-8') as outfile: json.dump(all_user_details, outfile, ensure_ascii=False, indent=2) print(f"已导出{len(all_user_details)}条用户数据到user_data.json") if __name__ == "__main__": with client: client.loop.run_until_complete(main(phone))
内容的提问来源于stack exchange,提问作者Rafael Sanches
相关产品推荐
相关产品推荐

