多线程运行Telegram客户端使用Telethon、Pyrogram收发消息问题求助
问题原因排查
1. Telethon初始版本问题
- 异步初始化函数
telethon_telegram_init未被await执行,直接赋值得到的是协程对象而非客户端实例,因此调用send_message触发属性错误 - 未启动asyncio事件循环,客户端的消息监听handler无法运行
- 多线程内直接调用异步方法
send_message未提交到客户端所在的事件循环执行,方法无法生效
2. Pyrogram版本问题
- Pyrogram v2.x版本中
send_message为异步方法,同步调用未执行也未抛出异常(未被await的协程会被直接回收) - 客户端启动后未运行事件循环,消息监听handler无法触发
3. asyncio版本问题
- 队列消费逻辑中使用同步
time.sleep(1)阻塞了整个asyncio事件循环,导致客户端事件处理、消息发送逻辑全部卡住,需替换为await asyncio.sleep(1) loop.run_in_executor提交了异步函数telethon_telegram_generate,同步线程中调用异步函数只会返回协程对象不会执行- 事件循环被阻塞后Telethon客户端的更新监听逻辑无法运行,因此无法接收消息
可行实现方案
推荐使用「asyncio事件循环 + 任务队列 + 跨线程提交任务」的方案,不要直接在多线程中操作Telegram客户端实例,所有客户端操作都提交到其所在的asyncio事件循环中执行:
import asyncio from telethon import TelegramClient, events from socketserver import ThreadingMixIn, TCPServer, StreamRequestHandler # 配置参数 API_ID = 替换为你的API_ID API_HASH = 替换为你的API_HASH PHONE_NUMBER = 替换为你的绑定手机号 RECIPIENT = 'me' # 全局客户端实例 client = None # TCP服务端实现 class ThreadingTCPServer(ThreadingMixIn, TCPServer): daemon_threads = True class TCPRequestHandler(StreamRequestHandler): def handle(self): msg = f"New connection from {self.client_address[0]}:{self.client_address[1]}" print(msg) # 跨线程将消息提交到asyncio队列 asyncio.run_coroutine_threadsafe(self.server.queue.put(msg), self.server.loop) async def run_tcp_server(loop, queue): server = ThreadingTCPServer(('127.0.0.1', 8080), TCPRequestHandler) server.loop = loop server.queue = queue # 同步的serve_forever放到线程池运行,避免阻塞事件循环 await loop.run_in_executor(None, server.serve_forever) # 客户端初始化 async def init_telegram(): global client client = TelegramClient('my_account', API_ID, API_HASH) await client.start(phone=PHONE_NUMBER) client.parse_mode = 'html' # 消息接收handler @client.on(events.NewMessage(incoming=True, func=lambda e: e.is_private)) async def on_private_msg(event): await event.respond('Thank you for your message') return client # 队列消费逻辑 async def consume_queue(queue): while True: msg = await queue.get() try: await client.send_message(RECIPIENT, msg) print(f"Sent message: {msg}") except Exception as e: print(f"Send failed: {str(e)}") finally: queue.task_done() await asyncio.sleep(1) async def main(): loop = asyncio.get_running_loop() msg_queue = asyncio.Queue() # 初始化Telegram客户端 await init_telegram() # 启动TCP服务 loop.create_task(run_tcp_server(loop, msg_queue)) # 启动队列消费 loop.create_task(consume_queue(msg_queue)) # 测试主动发消息 await msg_queue.put("Hello from initial call") # 运行Telegram客户端事件循环 await client.run_until_disconnected() if __name__ == "__main__": asyncio.run(main())
运行验证
- 首次启动会要求输入验证码,完成登录
- 执行
curl http://localhost:8080即可触发消息发送到你的Telegram账号 - 给你的账号发私信会自动收到回复
内容的提问来源于stack exchange,提问作者auino
相关产品推荐
相关产品推荐

