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

向独立线程运行的asyncio事件循环添加新任务无响应如何解决

问题原因

  • 获取事件循环逻辑错误:你在主线程中调用asyncio.get_event_loop()获取到的是主线程对应的事件循环,而不是你在子线程CraftSubscriptionThread中运行的事件循环,后续调用asyncio.run_coroutine_threadsafe时把协程提交到了没有启动运行的主线程循环,自然不会执行。
  • 事件循环生命周期错误:你用asyncio.run启动子线程循环,这个API会在传入的协程执行完成后直接关闭事件循环,就算你后续拿到了正确的子线程循环,也无法再提交新任务。
  • WebSocket连接复用逻辑错误:你在每个订阅任务中都用async with self._ws_client新建连接,上下文退出时会主动关闭连接,会直接打断正在运行的其他订阅任务,也无法支撑多订阅复用连接的需求。

修复方案

你需要主动存储子线程的事件循环实例,修改循环启动逻辑让它长期运行,同时复用同一个WebSocket连接处理多个订阅:

import asyncio
import threading
from typing import Union, Callable

class Craft:
    # 新增属性存储循环、线程、复用的WebSocket会话实例
    _subscription_loop: asyncio.AbstractEventLoop = None
    _subscription_thread: threading.Thread = None
    _ws_session = None

    async def exec_subscription(self, session, subscription: str, variable_values: dict, callback: Callable) -> None:
        """Execute a subscription on the GraphQL API."""
        async for response in session.subscribe(gql(subscription), variable_values=variable_values):
            callback(response)

    def subscribe_job(self, job_id: int, callback: Callable) -> Union[bool, None]:
        """Subscribe to a job, receive the job information every time it is updated and pass it to a callback function."""
        # 构建订阅语句和参数的逻辑保持不变
        subscribe_job = """
            subscription jobSubscription($jobId: Int!) {
                jobSubscriptionById(id: $jobId) {
                    job {...}
                }
            }
        """
        variable_values = {
            'jobId': job_id
        }

        # 循环已初始化的场景直接提交新任务
        if self._subscription_loop is not None and self._subscription_thread.is_alive():
            asyncio.run_coroutine_threadsafe(
                self.exec_subscription(self._ws_session, subscribe_job, variable_values, callback),
                self._subscription_loop
            )
            return True

        # 未初始化场景启动子线程、循环和WebSocket连接
        def _run_loop():
            self._subscription_loop = asyncio.new_event_loop()
            asyncio.set_event_loop(self._subscription_loop)
            
            async def _init_ws_and_run():
                async with self._ws_client as session:
                    self._ws_session = session
                    # 启动第一个订阅任务
                    self._subscription_loop.create_task(
                        self.exec_subscription(session, subscribe_job, variable_values, callback)
                    )
                    # 长期阻塞保持循环、连接存活
                    while True:
                        await asyncio.sleep(3600)
            
            self._subscription_loop.run_until_complete(_init_ws_and_run())

        self._subscription_thread = threading.Thread(
            name='CraftSubscriptionThread',
            daemon=True,
            target=_run_loop
        )
        self._subscription_thread.start()
        return True

额外优化建议

  • 可添加线程锁避免多线程同时调用subscribe_job时重复初始化循环
  • 可增加异常捕获逻辑,在WebSocket断开时自动重连、恢复所有订阅任务

内容的提问来源于stack exchange,提问作者Alexis.Rolland

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 12:18:02