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

如何在SocketIO事件处理器外部调用客户端?解决协程未等待警告

问题:在Socket.IO事件处理器外部调用emit/call方法

我尝试将sio实例传递给另一个类,在其中向客户端发送若干call请求,完成后返回process_msg函数。但由于AsyncServer.call是异步函数,这种方式无法正常运行,仅收到警告:RuntimeWarning: coroutine 'AsyncServer.call' was never awaited。我想知道是否可以在事件处理器外部emit或call客户端?该如何实现?

初始问题代码:

import socketio
from aiohttp import web


class TestClass:
    def __init__(self):
        self.sio = socketio.AsyncServer(async_mode='aiohttp')
        self.response_functions()
        self.app = web.Application()
        self.sio.attach(self.app)
        web.run_app(self.app, host="0.0.0.0", port=5000, print=None, access_log=None)
        
    def response_functions(self):

        @self.sio.on("message")
        async def process_msg(sid, msg):
            ac = AnotherClass(self.sio, sid)
            ac.do_some_action_here()
            await self.sio.emit("my_event", "is_finished")

class AnotherClass:
    def __init__(self, sio, sid):
        self.sio = sio
        self.sid = sid

    def do_some_action_here(self):
        # send some calls to the client, and wait to finish
        self.sio.call("foo", "bar1", to=self.sid)
        self.sio.call("foo", "bar2", to=self.sid)
        self.sio.call("foo", "bar3", to=self.sid)
        # return back to the event handler

TestClass()

编辑1:异步改造可行但不符合实际需求

将方法改为异步并使用await可以正常运行:

async def process_msg(sid, msg):
    ac = AnotherClass(self.sio, sid)
    await ac.do_some_action_here()
    await self.sio.emit("my_event", "is_finished")

async def do_some_action_here(self):
    # send some calls to the client, and wait to finish
    await self.sio.call("foo", "bar1", to=self.sid)

但实际代码中存在大量同步函数和类,这种全异步改造的方式不可行。我希望实现一个单例线程辅助类,在代码各处都能调用它向客户端发送事件,且无需等待事件完全完成。


编辑2:最终解决方案:线程+事件循环

通过将辅助类运行在独立线程中,并获取Socket.IO服务器的事件循环,即可在同步代码中调用异步的call方法,同时不阻塞主任务。以下是最终示例代码:

import socketio
from aiohttp import web
import asyncio
from threading import Thread

class TestClass:
    def __init__(self):
        self.sio = socketio.AsyncServer(async_mode='aiohttp')
        self.response_functions()
        self.app = web.Application()
        self.sio.attach(self.app)
        web.run_app(self.app, host="0.0.0.0", port=5000, print=None, access_log=None)
        
    def response_functions(self):
        @self.sio.on("message")
        async def process_msg(sid, msg):
            server_loop = asyncio.get_event_loop()
            AnotherClass(self.sio, sid, server_loop).start()

class AnotherClass(Thread):
    def __init__(self, sio, sid, server_loop):
        super(AnotherClass, self).__init__(daemon=True)
        self.sio = sio
        self.sid = sid
        self.server_loop = server_loop

    def run(self):
        # Do some calls and wait for each call to finish
        self.handler("foo", "bar1")
        self.handler("foo", "bar2")
        self.handler("foo", "bar3")

    def handler(self, event, msg):
        async def send(event, msg):
            await self.sio.call(event, msg, to=self.sid)
        asyncio.run_coroutine_threadsafe(send(event, msg), self.server_loop).result()

TestClass()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 23:02:50