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

如何优化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.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 01:40:08