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

