如何让基于NATS的WebSocket服务器支持多客户端同时连接?
问题描述
运行提供的WebSocket服务器代码时,第一个客户端能正常连接并接收NATS消息,但第二个客户端连接时触发超时错误,无法建立连接。
问题根源
- 事件循环阻塞:在
hello函数中调用SubscribeHandler.execute()时,该方法会新建独立事件循环并执行loop.run_forever(),直接阻塞当前处理WebSocket连接的协程,导致服务器无法响应后续的客户端连接请求。 - 重复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())
测试验证
- 启动
server.py,查看NATS连接成功和WebSocket服务器启动的日志。 - 启动多个
client.py实例,所有客户端均可正常建立连接。 - 执行
nats pub hello.world "test message",所有在线客户端都会收到这条消息。
内容的提问来源于stack exchange,提问作者Takewood
相关产品推荐
相关产品推荐

