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

Telethon如何并行收发消息?如何实现事件与非事件函数并行运行?

问题原因

代码中使用同步阻塞的input()函数是核心诱因。asyncio基于单线程事件循环调度所有异步任务,input()执行时会阻塞整个线程,导致事件循环无法调度Telethon的消息接收回调,只有用户完成输入、enviar函数逻辑继续执行后,事件循环才有机会处理收到的新消息,也就出现了收消息逻辑只有等发消息逻辑执行完才运行的现象。

解决方案

1. 替换同步输入为异步输入

使用异步输入工具替代原生同步input,避免阻塞事件循环,这里推荐使用aioconsole库的异步输入方法:

  • 先安装依赖:pip install aioconsole
  • 导入依赖:在代码开头添加from aioconsole import ainput
  • 替换enviar函数中的输入行:
    把
msg = input("Enviar: ")

替换为

msg = await ainput("Enviar: ")

2. 修复异常捕获的错误写法

你现有代码中except("Cannot send requests while disconnected")的写法不符合Python异常捕获规范,字符串不能作为异常捕获类型,建议修改为:

except Exception as e:
    if "Cannot send requests while disconnected" in str(e):
        print("Adeus!")
    else:
        print("\n======ERRO======\n")
        print(e)

3. 可选优化点

如果要进一步提升异步效率,可以把同步的文件读写操作替换为异步实现,使用aiofile库操作文件,避免小概率的IO阻塞影响事件循环调度。

修改后可正常并行运行的完整代码
from variables import api_id, api_hash
from telethon import TelegramClient, events, utils
import asyncio
import logging
from aioconsole import ainput

logging.basicConfig(format='[%(levelname) 5s/%(asctime)s] %(name)s: %(message)s',
                    level=logging.WARNING)

client = TelegramClient('anon', api_id, api_hash)

lock = asyncio.Lock()

async def info_me():
    me = await client.get_me()
    return me

async def enviar():
    while True:
        try:
            await lock.acquire()
            try:
                with open("msg.txt", "r") as f:
                    updates = f.read()
            finally:
                lock.release()
            
            print("==================")
            print(updates)
            print("==================")
            msg = await ainput("Enviar: ")
            if msg != "//":
                msg = msg.split(" | ")
                await client.send_message(msg[0], msg[1])
            else:
                await asyncio.sleep(0.1)
        except KeyboardInterrupt:
            print("Adeus!")
            break
        except Exception as e:
            print("\n======ERRO======\n")
            print(e)
            break
        
async def receber(event):
    try:
        sender = await event.get_sender()
        name = utils.get_display_name(sender)
        message = name + "::::::" + event.text + "\n"
        await lock.acquire()
        try:
            with open('msg.txt', 'a+') as f:
                f.write(message)
        finally:
            lock.release()
    except KeyboardInterrupt:
        print("Adeus!")
    except Exception as e:
        if "Cannot send requests while disconnected" in str(e):
            print("Adeus!")
        else:
            print("\n======ERRO======\n")
            print(e)
    
client.add_event_handler(receber, events.NewMessage)
async def main():
    me = await info_me()
    enviar_task = asyncio.create_task(enviar())
    await enviar_task
        
with client:
    client.loop.run_until_complete(main())

内容的提问来源于stack exchange,提问作者Joao Pedro Lourenco Affonso

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 22:24:02