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

Telethon跨客户端媒体下载问题及架构优化咨询

问题描述

我们尝试使用一个Telegram客户端持续从频道列表中流式获取消息并发送至Kafka,再由另一个Telegram客户端消费消息,通过client.download_media()下载关联的图片/视频媒体。但仅当两个客户端为同一账号时可行,不同账号则失效,不确定是否与session文件、access hash或其他因素有关。

核心诉求是解决异步媒体下载可能产生的大量积压问题,以及避免服务器宕机导致积压丢失,因此选用Kafka做短期存储,同时希望获取更优方案。

生产者端代码

async with client:
    messages = client.iter_messages(channel_id, limit=10)
    async for message in messages:
        print(message)
        if message.media is not None:
            # orig_media = message.media
            # converted_media = BinaryReader(bytes(orig_media)).tgread_object()
            # print('orig, media', orig_media)
            # print('converted media', converted_media)
            message_bytes = bytes(message) #convert to bytes
            producer.produce(topic, message_bytes)

不同客户端的消费者端代码

with self._client:
    #telethon.errors.rpcerrorlist.FileReferenceExpiredError: The file reference has expired and is no longer valid or it belongs to self-destructing media and cannot be resent (caused by GetFileRequest)
    try:
        self._client.loop.run_until_complete(self._client.download_media(orig_media, in_memory))
    except Exception as e:
        print(e)
问题分析与解决方案

场景支持性说明

不同账号跨客户端复用消息对象下载媒体的场景不被Telegram API支持,核心原因是:

  • Telegram媒体文件的file_reference(文件引用)与请求账号绑定且有有效期,生产者账号获取到的引用无法被其他账号复用。
  • 直接序列化整个Message对象传输时,其中的媒体元数据包含的是生产者账号的专属文件引用,消费者账号调用API时会因引用无效触发FileReferenceExpiredError。

针对性解决方案

1. 修改生产者:存储跨账号可用的媒体标识

不要序列化整个Message对象,而是提取能被任意有权限访问频道的账号复用的媒体核心信息:

import json

async with client:
    messages = client.iter_messages(channel_id, limit=10)
    async for message in messages:
        print(message)
        if message.media is not None:
            # 提取跨账号可用的媒体关键信息
            media_info = {
                "channel_id": channel_id,
                "message_id": message.id,
                "media_id": message.media.id,
                "media_access_hash": message.media.access_hash,
                "file_reference": message.media.file_reference.hex() if hasattr(message.media, 'file_reference') else None
            }
            # 转为JSON字符串发送到Kafka
            producer.produce(topic, json.dumps(media_info).encode('utf-8'))

2. 修改消费者:重新获取有效媒体引用

消费者拿到媒体标识后,通过重新拉取消息或构造合法的媒体输入对象来获取有效文件引用:

import json
from telethon.tl.types import InputPhoto, InputDocument

# 假设从Kafka消费到的是media_info的JSON字符串
media_info = json.loads(consumed_message.value.decode('utf-8'))

with self._client:
    try:
        # 方案1:重新拉取整条消息,直接获取有效媒体信息
        message = self._client.loop.run_until_complete(
            self._client.get_messages(media_info['channel_id'], ids=media_info['message_id'])
        )
        self._client.loop.run_until_complete(self._client.download_media(message.media, in_memory))
        
        # 方案2:用媒体ID、access_hash和file_reference构造输入对象(适合无需完整消息的场景)
        # if isinstance(message.media, InputPhoto):
        #     input_media = InputPhoto(
        #         media_info['media_id'], 
        #         media_info['media_access_hash'], 
        #         bytes.fromhex(media_info['file_reference'])
        #     )
        # elif isinstance(message.media, InputDocument):
        #     input_media = InputDocument(
        #         media_info['media_id'], 
        #         media_info['media_access_hash'], 
        #         bytes.fromhex(media_info['file_reference'])
        #     )
        # self._client.loop.run_until_complete(self._client.download_media(input_media, in_memory))
    except Exception as e:
        print(e)

3. 积压与可靠性优化方案

  • Kafka配置优化:给消息设置合理的过期时间(比如24小时),避免存储已失效的file_reference;开启消息确认机制,确保生产者消息写入Kafka后再标记Telegram消息已处理,防止丢失。
  • 消费者重试机制:遇到FileReferenceExpiredError时,自动触发重新拉取消息的逻辑,获取最新的有效文件引用。
  • 减少中转压力:如果目标频道是公开可订阅的,消费者可以直接订阅频道实时接收消息,无需通过Kafka中转,从根源上避免积压和引用过期问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 15:40:37