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

Rust异步WebSocket锁竞争问题:主线程无法获取锁

问题原因分析

核心问题在于异步锁的持有时间过长,导致主线程永远无法获取锁:

  1. 在start_websocket的tokio任务中,用try_lock拿到锁后,ws_locked.try_next().await会异步等待消息,而tokio的Mutex在await期间会一直持有锁(异步锁不会在挂起时自动释放)。
  2. 紧接着你把ws_locked传给handle_message,意味着整个handle_message执行期间锁都被持有,进一步拉长了锁的占用时间。
  3. 消息循环是无限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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 11:52:34