Rust异步WebSocket锁竞争问题:主线程无法获取锁
问题原因分析
核心问题在于异步锁的持有时间过长,导致主线程永远无法获取锁:
- 在
start_websocket的tokio任务中,用try_lock拿到锁后,ws_locked.try_next().await会异步等待消息,而tokio的Mutex在await期间会一直持有锁(异步锁不会在挂起时自动释放)。 - 紧接着你把
ws_locked传给handle_message,意味着整个handle_message执行期间锁都被持有,进一步拉长了锁的占用时间。 - 消息循环是无限loop,只要WebSocket连接正常,锁会被持续持有,主线程的
try_lock自然一直失败。
另外,send_data里用try_lock加轮询重试的方式并不适合异步场景,tokio的Mutex提供了lock().await方法,可以异步等待锁释放,比轮询高效得多。
修复方案
1. 缩小锁的持有范围
在消息处理线程中,只在调用try_next时持有锁,拿到消息后立即释放锁,再处理消息:
- 不要把锁传递给
handle_message,如果handle_message需要发送消息,单独获取锁即可。
2. 使用异步锁的正确姿势
放弃try_lock,改用lock().await来异步等待锁,避免轮询浪费资源。
3. 优化send_data逻辑
直接用lock().await获取锁,不需要手动轮询重试。
修复后的代码示例
调整start_websocket方法
pub async fn start_websocket( &mut self, auth: &Auth, ) -> Result<(), Box<dyn std::error::Error>> { let (ws_original, _) = connect_async(WEBSOCKET_ENDPOINT).await?; debug!("Connected to websocket"); self.websocket = Some(Arc::new(Mutex::new(ws_original))); let callback_clone = self.callback.clone().unwrap(); let websocket_clone = self.websocket.clone().unwrap(); let ws_thread = tokio::spawn(async move { loop { debug!("Waiting for message"); // 异步等待获取锁 let mut ws_locked = match websocket_clone.lock().await { Ok(ws) => ws, Err(_) => { debug!("Websocket lock poisoned"); tokio::time::sleep(Duration::from_millis(100)).await; continue; } }; // 只在try_next时持有锁,拿到消息后立即释放 let message = match ws_locked.try_next().await { Ok(Some(msg)) => msg, Ok(None) => { debug!("No more messages"); break; // 连接关闭,退出循环 } Err(e) => { debug!("Error while receiving message: {:?}", e); tokio::time::sleep(Duration::from_millis(100)).await; continue; } }; // 锁已经自动释放,现在处理消息 debug!("Message received: {:?}", message); handle_message(message, &callback_clone).await; debug!("Message handled"); } }); self.ws_thread = Some(ws_thread); Ok(()) }
调整send_data方法
async fn send_data(&mut self, data: serde_json::Value) -> bool { let ws = match self.websocket.as_ref() { Some(ws) => ws, None => { info!("Websocket not connected"); return false; } }; // 异步等待获取锁 let mut ws_locked = match ws.lock().await { Ok(ws) => ws, Err(_) => { info!("Websocket lock poisoned"); return false; } }; let send_result = ws_locked .send(tokio_tungstenite::tungstenite::Message::Text( data.to_string(), )) .await; match send_result { Ok(_) => { info!("Data sent successfully"); true } Err(e) => { info!("Failed to send data: {:?}", e); false } } }
调整handle_message(如果需要发送消息)
如果handle_message需要向WebSocket发送数据,不要传递锁,而是在内部单独获取:
async fn handle_message( message: tokio_tungstenite::tungstenite::Message, callback: &Arc<Mutex<dyn WebSocketCallback + Send>>, // 不再传递ws_locked ) { // 处理消息逻辑... // 如果需要发送消息: // let ws = get_websocket_instance(); // 假设能拿到WebSocket的Arc<Mutex>引用 // let mut ws_locked = ws.lock().await.unwrap(); // ws_locked.send(...).await; }
额外设计建议
- 避免在长时间运行的异步操作中持有锁,这会导致其他任务无法访问共享资源,引发阻塞或饥饿问题。
- tokio的
Mutex是为异步场景设计的,不要用try_lock做轮询,lock().await是更高效的异步等待方式。 - 如果
WebSocketStream本身是线程安全的(Send + Sync),可以考虑不需要Mutex,直接用Arc共享,但tungstenite的WebSocketStream通常不是Sync的,所以还是需要Mutex包裹。
内容的提问来源于stack exchange,提问作者P0W
相关产品推荐
相关产品推荐

