在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
相关产品推荐
相关产品推荐

