如何优化Telegram数据下载转存CSV的异步代码性能
优化Telegram异步数据抓取函数的问题解答
你的原始代码如下:
async def async_update_database(self, source_id: str, source_hash: str, message_id_1: str, message_id_n: str, file: str): entity = types.InputPeerChannel(int(source_id), int(source_hash)) n_messages = 0 n_comments = 0 with open(file, 'a', encoding = "utf-8") as fout: for message_id in range(int(message_id_1), int(message_id_n) + 1): message = await self.client.get_messages(entity, ids = message_id) if message == None or message.message == "" or message.replies == None or message.replies.replies == 0: continue fout.write("M," + str(int(message.date.timestamp())) + ',' + str(message_id) + ',' + repr(message.message) + '\n') n_messages += 1 async for comment in self.client.iter_messages(entity, reply_to = message_id): if isinstance(comment.sender, types.User) and not comment.sender.bot: fout.write("C," + str(int(comment.date.timestamp())) + ',' + str(comment.sender.id) + ',' + repr(comment.message) + '\n') n_comments += 1 return (str(n_messages) + " " + str(n_comments))
针对你的三个疑问,逐一解答:
1. 当前的文件追加写入方式是否高效?
不高效。现在每获取一条数据就调用一次fout.write(),频繁的磁盘IO会严重拖慢整体速度,尤其是要写入千万级行数据时,这个影响会被放大。另外用repr()处理消息内容适配CSV格式也有隐患——如果消息里包含引号、逗号或换行符,很容易导致CSV格式错乱。
优化建议:
- 启用文件缓冲区:打开文件时指定
buffering参数,比如buffering=1024*1024(1MB缓冲区),让系统攒够一定量数据再一次性写入磁盘,减少IO次数。 - 用
csv模块替代手动拼接:Python标准库的csv.writer会自动处理特殊字符,同时批量写入的效率更高。可以先把要写入的行收集到列表中,每隔1000行左右批量写入一次,最后再处理剩余内容。 - 减少文件开关次数:如果外部频繁调用该函数,可将文件句柄管理移到函数外部,或保持文件打开状态,避免重复打开关闭的开销。
2. 异步调用是否存在问题?
存在串行阻塞的核心问题。当前代码在for循环里逐个await get_messages(),相当于每次只能获取一条消息,完全没利用异步IO的并发优势。处理评论时也是串行遍历,整体效率和同步代码差异不大。
优化建议:
- 批量并发获取消息:把
message_id分成若干批次(比如每批次20个id),用asyncio.gather()同时发起多个get_messages请求,注意并发数控制在10-30之间(Telegram有API速率限制)。 - 并发处理评论抓取:将每个消息的评论抓取包装成异步任务,加入并发队列,但要注意文件写入的线程安全——异步环境下多任务同时写文件会导致数据错乱,建议每个任务先把数据收集到内存,最后统一写入,或用
asyncio.Lock()保护文件写入操作。 - 消除不必要的串行等待:确保所有可并发的操作都采用并发方式执行,减少等待时间。
3. Telethon库函数的使用是否正确?
整体逻辑没问题,但有两处关键优化点:
get_messages支持批量获取:当前每次只传一个message_id,而该函数可接收id列表(如ids=[1,2,3]),一次API调用就能获取多条消息,大幅减少请求次数,提升效率。- 消息判断逻辑可简化:Telethon的
get_messages找不到对应id时会返回None,用if not message可同时覆盖None和空消息的情况;另外message.replies不存在时直接访问message.replies.replies会抛出AttributeError,需先判断message.replies是否存在。
优化后的Telethon调用示例:
# 批量获取消息示例 message_ids = list(range(int(message_id_1), int(message_id_n)+1)) # 每20个id分为一组 chunks = [message_ids[i:i+20] for i in range(0, len(message_ids), 20)] for chunk in chunks: messages = await self.client.get_messages(entity, ids=chunk) for message in messages: if not message or not message.message or not message.replies or message.replies.replies == 0: continue # 后续消息处理逻辑...
另外,iter_messages的使用是正确的,指定reply_to=message_id可正确遍历该消息的所有评论,且会自动处理分页,无需手动干预。
内容的提问来源于stack exchange,提问作者I.S.M.
相关产品推荐
相关产品推荐

