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

如何在同一脚本中集成Telethon与Pika实现双向消息流转

最优方案选型:Asyncio全异步架构

选型对比

不推荐线程/多进程的核心原因

  • Telethon本身是原生异步库,强制放到同步线程中运行需要额外做线程安全封装,极易出现事件循环冲突、会话锁死问题
  • 多进程需要额外处理Telegram会话共享、RabbitMQ连接复用逻辑,大幅提升代码复杂度,IO密集型场景下没有任何性能优势
  • 线程场景下Pika的同步监听会阻塞Telethon的事件循环,反之亦然,需要额外加锁控制共享资源,故障率高

Asyncio方案的优势

  • 和Telethon的异步底层架构天然兼容,不需要额外做适配封装
  • 两个业务逻辑共用同一个事件循环,无需跨线程/进程通信,资源开销极低,运行稳定性更高
  • 异步生态完善,RabbitMQ有原生异步客户端可以直接复用

具体实现思路

推荐优先选择全异步适配方案,有两种实现路径可选:

  1. 替换Pika为全异步RabbitMQ客户端aio-pika,和Telethon的异步生态完全匹配,开发成本最低
  2. 如果必须保留Pika依赖,使用pika.adapters.asyncio_connection.AsyncioConnection适配器,将Pika的监听逻辑挂载到Telethon的同一个事件循环中运行

核心代码示例(aio-pika方案)

import asyncio
from telethon import TelegramClient, events
import aio_pika

# 配置参数
API_ID = 替换为你的Telegram API ID
API_HASH = "替换为你的Telegram API HASH"
TG_SESSION_NAME = "tg_session"
RABBITMQ_ADDR = "amqp://guest:guest@localhost/"
QUEUE_TG_IN = "telegram_receive"
QUEUE_TG_OUT = "telegram_send"
TARGET_CHAT_ID = 替换为消息发送目标的Telegram聊天ID

async def main():
    # 初始化Telegram客户端并登录
    tg_client = TelegramClient(TG_SESSION_NAME, API_ID, API_HASH)
    await tg_client.start()
    
    # 初始化RabbitMQ持久化连接
    rabbit_conn = await aio_pika.connect_robust(RABBITMQ_ADDR)
    async with rabbit_conn:
        channel = await rabbit_conn.channel()
        # 声明两个持久化队列
        await channel.declare_queue(QUEUE_TG_IN, durable=True)
        send_queue = await channel.declare_queue(QUEUE_TG_OUT, durable=True)
        
        # 注册Telegram消息回调:收到消息推送到RabbitMQ队列
        @tg_client.on(events.NewMessage())
        async def tg_msg_handler(event):
            await channel.default_exchange.publish(
                aio_pika.Message(body=event.text.encode(), delivery_mode=aio_pika.DeliveryMode.PERSISTENT),
                routing_key=QUEUE_TG_IN
            )
        
        # 注册RabbitMQ消费回调:收到消息推送到Telegram
        async def rabbit_msg_handler(message: aio_pika.IncomingMessage):
            async with message.process():
                send_text = message.body.decode()
                await tg_client.send_message(TARGET_CHAT_ID, send_text)
        
        # 启动RabbitMQ消费者
        await send_queue.consume(rabbit_msg_handler)
        # 启动Telegram客户端监听,直到手动断开
        await tg_client.run_until_disconnected()

if __name__ == "__main__":
    asyncio.run(main())

关键注意事项

  • 所有IO操作必须使用异步接口,不要在异步回调中调用任何同步阻塞代码,否则会卡住整个事件循环
  • RabbitMQ连接使用connect_robust方法,可以自动处理网络波动导致的断连重连问题
  • Telethon的会话文件不要多进程共享,避免出现文件锁冲突

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 13:06:03