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

Tornado WebSocket客户端如何异步实现on_message?(协程未被等待警告)

解决Tornado WebSocketClient中on_message异步执行的问题

你遇到的RuntimeWarning根源很明确:Tornado的websocket_connect方法接收的on_message_callback参数默认期望一个同步函数,当你传入异步协程(带async的on_message)时,Tornado只会直接调用这个协程对象,但不会自动await它,导致协程从未被执行,从而触发警告。

下面给你两种可行的解决方案,结合你的代码来修改:

方案一:放弃回调,改用主动循环读取消息(推荐)

这种方式更符合Tornado异步编程的习惯,逻辑也更清晰。我们在连接建立后,启动一个循环主动读取WebSocket消息,然后手动await你的异步on_message函数:

import tornado.websocket
from tornado.queues import Queue
from tornado import gen
import json
import tornado.ioloop

q = Queue()

class WebsocketClient():
    def __init__(self, url, connections):
        self.url = url
        self.connections = connections
        print("CLIENT started")
        print("CLIENT initial connections: ", len(self.connections))

    async def send_message(self):
        async for message in q:
            try:
                msg = json.loads(message)
                print(message)
                await gen.sleep(0.001)
            finally:
                q.task_done()

    async def update_connections(self, connections):
        self.connections = connections
        print("CLIENT updated connections: ", len(self.connections))

    async def on_message(self, message):
        await q.put(message)
        await gen.sleep(0.001)

    async def _message_loop(self):
        # 持续读取WebSocket消息
        while True:
            message = await self.client.read_message()
            if message is None:
                # 消息为None表示连接已关闭,退出循环
                print("WebSocket connection closed")
                break
            # 异步执行on_message
            await self.on_message(message)

    async def connect(self):
        # 建立连接,不使用on_message_callback
        self.client = await tornado.websocket.websocket_connect(url=self.url)
        print("Connected to WebSocket server")
        # 启动消息读取循环
        await self._message_loop()

# 示例运行代码
if __name__ == "__main__":
    async def main():
        client = WebsocketClient("ws://your-websocket-url", [])
        await client.connect()
    tornado.ioloop.IOLoop.current().run_sync(main)

方案二:保留回调,手动调度异步协程

如果你一定要使用on_message_callback,可以把异步的on_message包装在一个同步回调里,然后通过Tornado的IOLoop来调度执行这个协程:

# 只修改connect方法,其他代码不变
async def connect(self):
    def message_callback(message):
        # 将异步协程交给IOLoop调度,确保被await执行
        tornado.ioloop.IOLoop.current().add_callback(self.on_message, message)
    
    client = await tornado.websocket.websocket_connect(
        url=self.url, 
        on_message_callback=message_callback
    )

这个方案中,同步的message_callback被Tornado调用后,会把你的异步on_message加入IOLoop的任务队列,由Tornado负责执行并处理await逻辑,这样就不会再出现未被await的警告。

总结

你的思路没有根本性误区,只是没注意到Tornado对on_message_callback的同步要求。方案一更推荐,因为它的异步流程更直观,也更容易处理连接关闭等边界情况;方案二则适合需要保留回调模式的场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.07 19:42:31