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

多线程运行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())

运行验证

  1. 首次启动会要求输入验证码,完成登录
  2. 执行curl http://localhost:8080即可触发消息发送到你的Telegram账号
  3. 给你的账号发私信会自动收到回复

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 09:27:04