Telegram机器人多线程订阅机制异步IO报错求助
双线程Telegram机器人异步推送报错问题
我正在开发一个双线程Telegram机器人:一个线程负责生成数据,另一个线程管理机器人服务。用户发送/subscribe指令即可订阅数据更新,当数据生成线程完成数据生成后,会调用机器人线程将数据推送给所有订阅用户。目前功能大致可用,但推送数据时会抛出一系列异步IO相关报错,调整超时参数无效,推测是asyncio使用方式有误。
相关代码
import threading import asyncio from config import cfg from telegram import Update from telegram.ext import filters, ApplicationBuilder, ContextTypes, CommandHandler, MessageHandler import os subscribed = set() async def subscribe(update: Update, context: ContextTypes.DEFAULT_TYPE): subscribed.add(update.effective_chat.id) await context.bot.send_message( chat_id=update.effective_chat.id, text="You have subscribed" ) class TelegramService(): def __init__(self, control) -> None: self.control = control self.thread = threading.Thread(target=self.mainM) self.thread.start() def mainM(self): self.loop = asyncio.new_event_loop() asyncio.set_event_loop(self.loop) self.application = ApplicationBuilder().token(cfg.telegram_TOKEN).build() self.application.add_handler(CommandHandler('subscribe', subscribe)) self.application.run_polling() def sendOut(self, msg): asyncio.ensure_future(self.sendAsync(msg), loop=self.loop) async def sendAsync(self, msg): async with self.application.bot: for chat_id in subscribed: await self.application.bot.send_message(text = msg, chat_id=chat_id)
报错信息
获取更新时出错: httpx.ReadError: 任务异常未被捕获 future: <Task finished coro=<TelegramService.sendAsync() done, defined at .\communication\telegramService.py:83> exception=NetworkError('httpx.ReadError: ')> 回溯(最近的调用在最前面): 文件 "d:\Development\.Projects.Chast\ViGrabber2\.venv\lib\site-packages\httpcore\_exceptions.py", line 10, in map_exceptions yield 文件 "d:\Development\.Projects.Chast\ViGrabber2\.venv\lib\site-packages\httpcore\backends\asyncio.py", line 34, in read return await self._stream.receive(max_bytes=max_bytes) 文件 "d:\Development\.Projects.Chast\ViGrabber2\.venv\lib\site-packages\anyio\streams\tls.py", line 195, in receive data = await self._call_sslobject_method(self._ssl_object.read, max_bytes) 文件 "d:\Development\.Projects.Chast\ViGrabber2\.venv\lib\site-packages\anyio\streams\tls.py", line 137, in _call_sslobject_method data = await self.transport_stream.receive() 文件 "d:\Development\.Projects.Chast\ViGrabber2\.venv\lib\site-packages\anyio\_backends\_asyncio.py", line 1272, in receive raise ClosedResourceError from None anyio.ClosedResourceError 在处理上述异常期间,又发生了新的异常: ....
完整报错栈还包含多个httpcore.ReadError、httpx.ReadError,最终以RuntimeError: This HTTPXRequest is not initialized!和telegram.error.NetworkError结尾。
问题根源
- 跨线程提交异步任务的方式不安全:使用
asyncio.ensure_future在非事件循环线程提交任务,虽指定了loop,但未正确处理线程同步,导致bot的HTTP客户端资源被意外关闭。 - 误用
async with self.application.bot:application.bot在run_polling启动后已处于活跃状态,重复用async with会重新初始化客户端,引发连接冲突。 - 订阅集合非线程安全:直接用普通
set存储订阅ID,跨线程读写可能导致数据不一致。
修复方案
1. 线程安全提交异步任务
修改sendOut方法,用loop.call_soon_threadsafe确保跨线程提交任务的安全性:
def sendOut(self, msg): self.loop.call_soon_threadsafe( asyncio.create_task, self.sendAsync(msg) )
2. 移除不必要的bot上下文管理
sendAsync无需重新进入bot上下文,直接使用已初始化的实例:
async def sendAsync(self, msg): for chat_id in subscribed: try: await self.application.bot.send_message(text=msg, chat_id=chat_id) except Exception as e: print(f"推送消息给{chat_id}失败: {e}")
3. 线程安全的订阅集合
改用带锁的集合避免跨线程读写问题:
from threading import Lock subscribed = set() sub_lock = Lock() async def subscribe(update: Update, context: ContextTypes.DEFAULT_TYPE): chat_id = update.effective_chat.id with sub_lock: subscribed.add(chat_id) await context.bot.send_message( chat_id=chat_id, text="订阅成功" ) async def sendAsync(self, msg): with sub_lock: chat_ids = list(subscribed) # 复制一份避免遍历中集合被修改 for chat_id in chat_ids: try: await self.application.bot.send_message(text=msg, chat_id=chat_id) except Exception as e: print(f"推送消息给{chat_id}失败: {e}")
4. 推荐:改用单事件循环+多任务
无需用多线程,asyncio本身支持并发,将数据生成逻辑改为异步任务,与机器人服务共用一个事件循环,彻底避免跨线程异步问题:
import asyncio from config import cfg from telegram import Update from telegram.ext import filters, ApplicationBuilder, ContextTypes, CommandHandler, MessageHandler subscribed = set() sub_lock = asyncio.Lock() async def subscribe(update: Update, context: ContextTypes.DEFAULT_TYPE): chat_id = update.effective_chat.id async with sub_lock: subscribed.add(chat_id) await context.bot.send_message( chat_id=chat_id, text="订阅成功" ) async def data_generator(bot): # 替换为你的实际数据生成逻辑 while True: await asyncio.sleep(10) new_data = "新生成的数据内容" await send_all_subscribers(new_data, bot) async def send_all_subscribers(msg, bot): async with sub_lock: chat_ids = list(subscribed) for chat_id in chat_ids: try: await bot.send_message(text=msg, chat_id=chat_id) except Exception as e: print(f"推送消息给{chat_id}失败: {e}") async def main(): application = ApplicationBuilder().token(cfg.telegram_TOKEN).build() application.add_handler(CommandHandler('subscribe', subscribe)) # 启动数据生成任务 asyncio.create_task(data_generator(application.bot)) await application.run_polling() if __name__ == "__main__": asyncio.run(main())
内容的提问来源于stack exchange,提问作者Archeg
相关产品推荐
相关产品推荐

