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

