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

在Rust中使用Tungstenite实现WebSocket并发读写的问题

解决Rust Tungstenite WebSocket并发读写的锁阻塞问题

问题根源

核心问题是读循环长期持有WebSocket的MutexGuard,导致写操作无法获取锁而永久阻塞。异步环境下,Mutex的持有时间必须尽可能短,绝不能在整个循环周期内持续占用锁。

解决方案

1. 使用tokio异步Mutex

确保使用tokio::sync::Mutex而非标准库的std::sync::Mutex——异步场景下标准库Mutex会阻塞整个执行器,而tokio的Mutex专为异步设计,支持协程调度,不会卡死执行流程。

2. 缩小锁的作用域

修改UnifiedWebSocket的read和send方法,让每次读写操作仅在执行IO的瞬间持有锁,而非在外层循环中持续持有。

完整修改示例

定义UnifiedWebSocket枚举

use tokio::sync::Mutex;
use tungstenite::protocol::Message;
use tungstenite::{Error as WsError};
use tokio_tungstenite::WebSocketStream;
use std::sync::Arc;

// 统一普通与TLS类型的WebSocket内部实现
enum UnifiedWebSocketInner {
    Plain(WebSocketStream<tokio::net::TcpStream>),
    Tls(WebSocketStream<tokio_native_tls::TlsStream<tokio::net::TcpStream>>),
}

pub struct UnifiedWebSocket {
    inner: Arc<Mutex<UnifiedWebSocketInner>>,
}

impl UnifiedWebSocket {
    // 普通WebSocket构造示例(根据实际场景调整)
    pub async fn new_plain(stream: tokio::net::TcpStream) -> Self {
        let ws = tokio_tungstenite::connect_async(stream)
            .await
            .unwrap()
            .0;
        Self {
            inner: Arc::new(Mutex::new(UnifiedWebSocketInner::Plain(ws))),
        }
    }

    // TLS WebSocket构造示例(根据实际场景调整)
    pub async fn new_tls(stream: tokio::net::TcpStream, domain: &str) -> Self {
        let connector = tokio_native_tls::TlsConnector::new().unwrap();
        let tls_stream = connector.connect(domain, stream).await.unwrap();
        let ws = tokio_tungstenite::connect_async(tls_stream)
            .await
            .unwrap()
            .0;
        Self {
            inner: Arc::new(Mutex::new(UnifiedWebSocketInner::Tls(ws))),
        }
    }

    // 异步发送:仅在发送时临时持有锁
    pub async fn send(&self, msg: Message) -> Result<(), WsError> {
        let mut guard = self.inner.lock().await;
        match &mut *guard {
            UnifiedWebSocketInner::Plain(ws) => ws.send(msg).await,
            UnifiedWebSocketInner::Tls(ws) => ws.send(msg).await,
        }
    }

    // 异步读取:仅在读取时临时持有锁
    pub async fn read(&self) -> Result<Message, WsError> {
        let mut guard = self.inner.lock().await;
        match &mut *guard {
            UnifiedWebSocketInner::Plain(ws) => ws.next().await.unwrap_or(Err(WsError::ConnectionClosed)),
            UnifiedWebSocketInner::Tls(ws) => ws.next().await.unwrap_or(Err(WsError::ConnectionClosed)),
        }
    }
}

正确的读写循环实现

use tokio::task;

#[tokio::main]
async fn main() {
    // 建立连接并创建WebSocket实例
    let tcp_stream = tokio::net::TcpStream::connect("ws://your-websocket-server.com")
        .await
        .unwrap();
    let ws = UnifiedWebSocket::new_plain(tcp_stream).await;

    // 启动读循环:每次迭代临时获取锁,读完立即释放
    let ws_read_clone = ws.clone();
    task::spawn(async move {
        loop {
            match ws_read_clone.read().await {
                Ok(msg) => {
                    println!("Received message: {:?}", msg);
                    // 处理消息逻辑...
                }
                Err(e) => {
                    eprintln!("Read failed: {}", e);
                    break;
                }
            }
        }
    });

    // 启动写任务:随时可获取锁发送消息
    let ws_write_clone = ws.clone();
    task::spawn(async move {
        tokio::time::sleep(tokio::time::Duration::from_secs(2)).await;
        match ws_write_clone.send(Message::Text("Hello from async writer".into())).await {
            Ok(_) => println!("Message sent successfully"),
            Err(e) => eprintln!("Send failed: {}", e),
        }
    });

    // 阻塞主线程直到收到终止信号
    tokio::signal::ctrl_c().await.unwrap();
}

关键说明

  • 锁作用域最小化:read和send方法内部仅在执行实际IO操作时持有锁,操作完成后立即释放,让读写任务可以交替获取锁,避免互相阻塞。
  • 异步Mutex适配:tokio::sync::Mutex在等待锁时会主动让出执行权,不会阻塞整个异步运行时,这是异步并发场景的核心要求。
  • Arc共享实例:通过Arc克隆WebSocket的inner字段,让读写任务各自持有安全的引用,实现并发环境下的共享访问。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 23:37:39