tokio::mpsc::unbounded_channel消息丢失:多客户端Websocket服务器异常
问题分析与解决方案
核心问题分析
- 阻塞式睡眠破坏Tokio调度:广播线程使用
std::thread::sleep,会直接占用Tokio工作线程并阻塞调度,导致客户端消息接收、广播转发等任务无法获得执行机会,多客户端场景下消息传递异常。 - Mutex锁持有时间过长:
send_clients函数在遍历客户端、序列化数据包的全流程中一直持有客户端列表的Mutex锁,导致客户端连接、断开等操作被长时间阻塞,影响消息转发效率。 - 代码逻辑错误:客户端消息处理中,
handle_message调用时传入的&mut clients未在当前作用域定义,属于编译级错误。 - 冗余可变状态:
State中的clients已通过Arc<Mutex>实现共享可变,广播线程和send_clients无需持有可变引用,多余的mut会引发生命周期问题。
修正后的代码
use futures_util::{future, StreamExt, TryStreamExt, SinkExt}; use serde::{Deserialize, Serialize}; use serde_json::Value; use std::{ collections::HashMap, net::SocketAddr, sync::{Arc, Mutex}, time::Instant, }; use tokio::{net::TcpListener, sync::mpsc::{UnboundedSender, self}, time::sleep}; use tungstenite::Message; use uuid::Uuid; mod packets; use packets::*; static REFRESH_TIME: u64 = 2000; #[derive(Debug, Clone)] pub struct Client { pub send: UnboundedSender<Message> } #[derive(Debug, Default, Clone)] pub struct State { pub clients: Arc<Mutex<HashMap<Uuid, Client>>> } #[tokio::main] async fn main() { let g_state = State::default(); let state_clone = g_state.clone(); tokio::task::spawn(async move { let state = state_clone; loop { send_clients(&state).await; sleep(std::time::Duration::from_millis(REFRESH_TIME)).await; } }); let stream = TcpListener::bind("0.0.0.0:25342").await.unwrap(); loop { let (socket, addr) = stream.accept().await.unwrap(); let state_clone = g_state.clone(); tokio::task::spawn(async move { let self_uuid = Uuid::new_v4(); let state = state_clone; let (tx, mut rx) = mpsc::unbounded_channel::<Message>(); { let mut clients = state.clients.lock().unwrap(); clients.insert( self_uuid, Client { send: tx }, ); } let websocket = tokio_tungstenite::accept_async(socket).await.unwrap(); let (mut outgoing, mut incoming) = websocket.split(); let incoming_handle = incoming.try_for_each(|msg| { handle_message(msg, &state, self_uuid); future::ok(()) }); let forward_handle = tokio::task::spawn(async move { while let Some(message) = rx.recv().await { if let Err(e) = outgoing.send(message).await { eprintln!("Failed to send message to client {}: {}", self_uuid, e); break; } } }); tokio::select! { res = incoming_handle => { if let Err(e) = res { eprintln!("Incoming message error for client {}: {}", self_uuid, e); } } _ = forward_handle => {} } { let mut clients = state.clients.lock().unwrap(); println!("[x] Client disconnected! @ {}", self_uuid); clients.remove(&self_uuid); } }); } } async fn send_clients(state: &State) { let senders: Vec<UnboundedSender<Message>> = { let clients = state.clients.lock().unwrap(); clients.values().map(|client| client.send.clone()).collect() }; let clients_snapshot = { let clients = state.clients.lock().unwrap(); get_clients_packet_from_clients(&clients) }; let packet = serde_json::to_string(&clients_snapshot).unwrap(); let message = Message::Text(packet); for sender in senders { if let Err(e) = sender.send(message.clone()) { eprintln!("Failed to send broadcast message: {}", e); } } } fn handle_message(msg: Message, state: &State, client_id: Uuid) { match msg { Message::Text(text) => println!("Received message from client {}: {}", client_id, text), Message::Binary(_) => println!("Received binary message from client {}", client_id), _ => {} } }
关键修改说明
- 替换为Tokio异步睡眠:用
tokio::time::sleep替代std::thread::sleep,让Tokio在睡眠期间调度其他任务,避免阻塞工作线程。 - 缩短锁持有时间:先复制客户端Sender列表、生成数据包快照,立即释放锁后再逐个发送消息,减少锁竞争对其他操作的影响。
- 修复未定义变量错误:修正
handle_message的调用参数,通过state访问客户端列表。 - 优化任务生命周期:用
tokio::select!管理消息处理和转发任务,确保客户端断开时正确清理资源;处理转发失败的情况,避免无限阻塞。 - 移除冗余可变引用:删除
State和send_clients中不必要的mut修饰符,符合Rust所有权规则。
内容的提问来源于stack exchange,提问作者Stopper
相关产品推荐
相关产品推荐

