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

使用Mastodon.py搭配Trio实现异步通知流的问题

Mastodon.py异步流配合Trio运行的问题解决

问题描述

使用Mastodon.py异步拉取Mastodon通知并通过Trio运行时,遇到以下问题:

  • 调用mastodon.stream_user()后,程序卡在该调用处,无法执行后续的logger.info("Started user stream"),仅输出"Starting user stream"
  • 添加run_async=True、reconnect_async=True参数后,流会立即退出,无法持续运行,也无法并行执行其他代码

解决方案

核心要点

  • 异步流的正确使用:设置run_async=True时,stream_user()返回异步上下文管理器,必须通过async with启动和管理,不能直接传入trio.run()
  • 异步监听器:若需在回调中执行异步操作,应继承AsyncStreamListener而非同步的StreamListener,避免阻塞事件循环
  • 保持程序运行:通过await trio.sleep_forever()或并行启动其他任务,防止主协程结束导致程序退出

修正后的代码

...
from mastodon import Mastodon, AsyncStreamListener  # 替换为异步监听器基类
import trio
...

async def main():
    ...
    access_token = os.getenv("MASTODON_ACCESS_TOKEN")
    api_base_url = os.getenv("MASTODON_API_BASE_URL")

    # 初始化客户端并登录
    mastodon = Mastodon(
        access_token=access_token,
        api_base_url=api_base_url
    )
    # 验证账号并输出用户名
    logger.info(f"Logged in as {mastodon.account_verify_credentials()['username']}")
    ...

    logger.info("Starting user stream")
    # 以异步模式启动用户流,用async with托管
    async with mastodon.stream_user(TheStreamListener(), run_async=True, reconnect_async=True) as stream:
        logger.info("Started user stream")
        try:
            # 保持主协程存活,可在此处用trio.start_soon()启动其他并行任务
            await trio.sleep_forever()
        except KeyboardInterrupt:
            logger.info("Stopping user stream")
            await stream.close()  # 异步关闭流
            logger.info("Stopped user stream")

class TheStreamListener(AsyncStreamListener):  # 继承异步监听器
    # 回调方法可定义为async,同步操作也兼容
    async def on_update(self, status):
        logger.info(f"Got update: {status['content']}")
        ...
    
    async def on_notification(self, notification):
        if notification['type'] == 'mention':
            logger.opt(colors=True).info(f"Got <blue>mention</blue> from {notification['account']['username']}")
        ...

if __name__ == "__main__":
    trio.run(main)

关键修改说明

  • 将main改为异步函数,由trio.run()直接调用
  • 用async with托管异步流,确保流的生命周期被正确管理
  • 替换StreamListener为AsyncStreamListener,支持异步回调逻辑
  • 用await trio.sleep_forever()维持主协程运行,若有其他任务,可通过trio.start_soon()并行启动
  • 关闭流时使用异步的await stream.close(),适配异步上下文管理器的要求

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 09:51:08