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

如何实现WebSocket连接全局化,接收多客户端消息并查看日志

问题

我搭建了一个WebSocket服务,客户端发送消息后服务器会将消息回传。使用两个浏览器测试时,两个客户端都能收到各自的消息,但服务器日志中无法看到两条消息。请问如何实现WebSocket连接全局化以查看所有客户端消息?是否需要创建另一个WebSocket来接收并转发这些消息?

前端代码

import { isPlatformBrowser } from '@angular/common';
import { Inject, Injectable } from '@angular/core';
import { PLATFORM_ID } from '@angular/core';

@Injectable({
  providedIn: 'root',
})
export class WebsocketService {
  messages: Array<any> = [];
  chatWebSocket: WebSocket | undefined = undefined;

  createSocket() {
    if (isPlatformBrowser(this.platformId)) {
      this.chatWebSocket = new WebSocket('ws://localhost:3000/ws');
    }
  }

  send(msg: string) {
    console.log(this.chatWebSocket);
    (this.chatWebSocket as WebSocket).send(msg);
  }

  recv() {
    if (!this.chatWebSocket) return;
    this.chatWebSocket.addEventListener('message', (e) => {
      this.messages.push(e.data);
    });
  }

  constructor(@Inject(PLATFORM_ID) public platformId: Object) {}
}

服务端代码

pub async fn chat_ws(
    ws: WebSocketUpgrade,
    user_agent: Option<TypedHeader<UserAgent>>,
) -> impl IntoResponse {
    if let Some(TypedHeader(user_agent)) = user_agent {
        println!("`{}` connected", user_agent.as_str());
    }

    ws.on_upgrade(handle_chat_socket)
}

pub async fn handle_chat_socket(mut socket: WebSocket) {
    loop {
        if let Some(msg) = socket.recv().await {
            if let Ok(msg) = msg {
                match msg {
                    Message::Text(a) => {
                        socket.send(Message::Text(String::from(a.clone()))).await;
                    },
                    _ => println!("Other"),
                }
            }
        }
    }
}
解决方案

不需要额外创建WebSocket服务,你只需要在服务端维护一个全局的客户端连接集合,就能追踪所有活跃连接并查看所有客户端消息,具体实现步骤如下:

  1. 用线程安全容器管理连接集合
    在Rust中,使用Arc<Mutex<Vec<WebSocket>>>来存储所有活跃的WebSocket连接,确保多异步任务下的安全访问。

  2. 修改服务端逻辑

    • 新客户端连接时,将其加入全局集合;
    • 收到消息时先打印日志,这样就能看到所有客户端的消息;
    • 可选:将消息广播给所有客户端(如果需要多客户端通信);
    • 客户端断开连接时,从集合中移除对应的连接,避免内存泄漏。

修改后的服务端示例代码

use std::sync::{Arc, Mutex};
use actix_web::{web, IntoResponse};
use actix_web_actors::ws;
use actix_web::http::header::TypedHeader;
use actix_web::http::header::UserAgent;

// 定义全局连接存储类型
type WsConnections = Arc<Mutex<Vec<ws::WebSocket>>>;

// 服务配置:初始化连接集合并注入
pub fn configure_services(cfg: &mut web::ServiceConfig) {
    let connections = Arc::new(Mutex::new(Vec::new()));
    cfg.app_data(connections.clone())
       .route("/ws", web::get().to(chat_ws));
}

pub async fn chat_ws(
    ws: ws::WebSocketUpgrade,
    user_agent: Option<TypedHeader<UserAgent>>,
    connections: web::Data<WsConnections>,
) -> impl IntoResponse {
    if let Some(TypedHeader(user_agent)) = user_agent {
        println!("`{}` connected", user_agent.as_str());
    }

    // 将连接集合传递给处理函数
    ws.on_upgrade(|socket| handle_chat_socket(socket, connections))
}

pub async fn handle_chat_socket(mut socket: ws::WebSocket, connections: web::Data<WsConnections>) {
    // 新连接加入全局集合
    {
        let mut conns = connections.lock().unwrap();
        conns.push(socket.clone());
    }

    loop {
        match socket.recv().await {
            Some(Ok(msg)) => {
                match msg {
                    ws::Message::Text(text) => {
                        // 打印所有客户端消息到服务器日志
                        println!("收到客户端消息: {}", text);
                        
                        // 回传给发送者(原有逻辑)
                        if let Err(e) = socket.send(ws::Message::Text(text.clone())).await {
                            eprintln!("回传消息失败: {}", e);
                        }
                        
                        // 可选:广播消息给所有其他客户端
                        let mut conns = connections.lock().unwrap();
                        for conn in conns.iter_mut() {
                            if *conn != socket {
                                if let Err(e) = conn.send(ws::Message::Text(text.clone())).await {
                                    eprintln!("广播消息失败: {}", e);
                                }
                            }
                        }
                    },
                    ws::Message::Binary(bin) => {
                        println!("收到二进制消息,长度: {}", bin.len());
                        if let Err(e) = socket.send(ws::Message::Binary(bin)).await {
                            eprintln!("回传二进制消息失败: {}", e);
                        }
                    },
                    _ => println!("收到其他类型消息"),
                }
            },
            Some(Err(e)) => {
                eprintln!("连接错误: {}", e);
                break;
            },
            None => {
                println!("客户端断开连接");
                break;
            }
        }
    }

    // 连接断开后从全局集合移除
    let mut conns = connections.lock().unwrap();
    conns.retain(|conn| *conn != socket);
}

注意事项

  • 必须使用线程安全的容器(Arc<Mutex>),因为每个WebSocket连接运行在独立的异步任务中,多任务访问集合需要同步;
  • 处理消息发送时的错误,避免单个连接的异常影响整个服务;
  • 前端代码无需修改,保持现有连接逻辑即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 08:09:16