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

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结尾。


问题根源

  1. 跨线程提交异步任务的方式不安全:使用asyncio.ensure_future在非事件循环线程提交任务,虽指定了loop,但未正确处理线程同步,导致bot的HTTP客户端资源被意外关闭。
  2. 误用async with self.application.bot:application.bot在run_polling启动后已处于活跃状态,重复用async with会重新初始化客户端,引发连接冲突。
  3. 订阅集合非线程安全:直接用普通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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 13:58:02