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

如何让基于NATS的WebSocket服务器支持多客户端同时连接?

问题描述

运行提供的WebSocket服务器代码时,第一个客户端能正常连接并接收NATS消息,但第二个客户端连接时触发超时错误,无法建立连接。

问题根源
  1. 事件循环阻塞:在hello函数中调用SubscribeHandler.execute()时,该方法会新建独立事件循环并执行loop.run_forever(),直接阻塞当前处理WebSocket连接的协程,导致服务器无法响应后续的客户端连接请求。
  2. 重复NATS订阅:每个客户端连接都会创建新的NATS客户端实例并订阅同一主题,既浪费资源,又会导致同一条消息被重复处理。
修复方案

核心思路:服务器启动时仅建立一次NATS连接并订阅主题,维护所有在线WebSocket客户端列表,当NATS收到消息时,将消息广播给所有在线客户端。

修改后的代码

1. subscribehandlernats.py

重构为异步类,复用主事件循环,不再创建独立循环:

import asyncio
import signal
from nats.aio.client import Client as NATS
from nats.aio.errors import ErrConnectionClosed, ErrTimeout, ErrNoServers

class NatsSubscriber:
    def __init__(self, subject_list, message_callback):
        self.subject_list = subject_list
        self.message_callback = message_callback
        self.nc = NATS()

    async def connect(self):
        try:
            await self.nc.connect("nats://localhost:4222")
            print(f"Connected to NATS at {self.nc.connected_url.netloc}...")
        except ErrNoServers as e:
            print(e)
            raise
        except Exception as e:
            print(e)
            raise

        # 注册信号处理
        loop = asyncio.get_running_loop()
        def signal_handler():
            if self.nc.is_closed:
                return
            print("Disconnecting from NATS...")
            loop.create_task(self.nc.close())

        for sig in ('SIGINT', 'SIGTERM'):
            loop.add_signal_handler(getattr(signal, sig), signal_handler)

        # 订阅主题
        for subject in self.subject_list:
            await self.nc.subscribe(subject, cb=self._handle_message)
            print(f"Subscribed to: {subject}")

    async def _handle_message(self, msg):
        subject = msg.subject
        reply = msg.reply
        data = msg.data.decode()
        print(f"Received a message on '{subject} {reply}': {data}")
        # 调用外部回调处理消息分发
        await self.message_callback(subject, reply, data)

    async def close(self):
        if not self.nc.is_closed:
            await self.nc.close()

2. server.py

维护客户端列表,实现消息广播,复用主事件循环处理NATS和WebSocket:

import asyncio
import websockets
from subscribehandlernats import NatsSubscriber
import nest_asyncio

# 存储所有在线WebSocket客户端
connected_clients = set()

async def broadcast_message(subject, reply, data):
    """将NATS消息广播给所有在线客户端"""
    message = f"Received a message on '{subject} {reply}': {data}"
    if connected_clients:
        await asyncio.gather(
            *[client.send(message) for client in connected_clients],
            return_exceptions=True
        )

async def handle_websocket(websocket):
    # 客户端连接时加入列表
    connected_clients.add(websocket)
    try:
        # 保持连接,直到客户端断开
        await websocket.wait_closed()
    finally:
        # 客户端断开时移除
        connected_clients.discard(websocket)

async def main():
    # 初始化NATS订阅器
    nats_subscriber = NatsSubscriber(["hello.world"], broadcast_message)
    await nats_subscriber.connect()

    # 启动WebSocket服务器
    async with websockets.serve(handle_websocket, "localhost", 8765, ping_interval=None):
        print("WebSocket server running on ws://localhost:8765")
        await asyncio.Future()  # 保持服务器运行
    finally:
        await nats_subscriber.close()

if __name__ == "__main__":
    nest_asyncio.apply()
    asyncio.run(main())

3. client.py(补充缺失的asyncio导入)

import asyncio
import websockets

async def hello():
    uri = "ws://localhost:8765"
    async with websockets.connect(uri, ping_interval=None) as websocket:
        while True:
            greeting = await websocket.recv()
            print(f"<<< {greeting}")

if __name__ == "__main__":
    asyncio.run(hello())
测试验证
  1. 启动server.py,查看NATS连接成功和WebSocket服务器启动的日志。
  2. 启动多个client.py实例,所有客户端均可正常建立连接。
  3. 执行nats pub hello.world "test message",所有在线客户端都会收到这条消息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 06:12:05