如何解决Rust循环中Arc<RwLock<HashMap>>写入失败问题?
问题:WebSocket循环中写入Arc<RwLock>触发panic
我编写了一段Rust代码,通过WebSocket循环接收消息,并将消息中的用户信息写入Arc<RwLock<HashMap>>(作为用户数据库)。代码如下:
#[derive(Debug)] struct ChatUser{ name: String, password: String, } static NEW_USER_ID: AtomicUsize = AtomicUsize::new(1); type UserDB = Arc<RwLock<HashMap<usize, ChatUser>>>; async fn web_socket(mut ws: WebSocket, State(state): State<UserDB>) { let new_id = NEW_USER_ID.fetch_add(1, std::sync::atomic::Ordering::Relaxed); while let Ok(Message::Text(message)) = ws.next().await.unwrap(){ let user_info: Value = serde_json::from_str(&message.trim()).unwrap(); let (name, password) = (user_info["name"].to_string(), user_info["password"].to_string()); state.write().await.insert(new_id, ChatUser { name: name, password: password }).expect("cant write"); println!("{:?}",state); } }
但在循环内执行state.write().await.insert时触发panic,提示‘cant write’,循环外执行则正常。想知道能否将值移出循环处理,或有其他正确写入state的方式?运行代码并发送消息到Socket时,报错栈如下:
thread 'tokio-runtime-worker' panicked at src/main.rs:41:93: cant write stack backtrace: 0: rust_begin_unwind at /rustc/3f5fd8dd41153bc5fdca9427e9e05be2c767ba23/library/std/src/panicking.rs:652:5 1: core::panicking::panic_fmt at /rustc/3f5fd8dd41153bc5fdca9427e9e05be2c767ba23/library/core/src/panicking.rs:72:14 2: core::panicking::panic_display at /rustc/3f5fd8dd41153bc5fdca9427e9e05be2c767ba23/library/core/src/panicking.rs:262:5 3: core::option::expect_failed at /rustc/3f5fd8dd41153bc5fdca9427e9e05be2c767ba23/library/core/src/option.rs:1995:5 4: core::option::Option<T>::expect [............]
修复方案
问题根源
HashMap::insert方法返回Option<T>:当插入的key不存在时返回None,当key已存在时返回之前存储的值。你代码里用了expect("cant write"),这意味着第二次用同一个new_id插入时,insert返回None,直接触发了panic。
看你的代码,new_id是在循环外初始化的——每个WebSocket连接只会生成一个id,循环内每次收到消息都用同一个id插入,第二次发送消息时key必然重复,导致panic。
具体修复方式
方式一:每次消息生成新用户ID(业务需要每条消息对应新用户)
把new_id的生成逻辑移到循环内部,每次收到消息都生成唯一的新ID:
async fn web_socket(mut ws: WebSocket, State(state): State<UserDB>) { while let Ok(Message::Text(message)) = ws.next().await.unwrap(){ // 每次消息生成新ID let new_id = NEW_USER_ID.fetch_add(1, std::sync::atomic::Ordering::Relaxed); let user_info: Value = serde_json::from_str(&message.trim()).unwrap(); let (name, password) = (user_info["name"].to_string(), user_info["password"].to_string()); // 不需要expect,因为每次都是新ID,不会重复 let _old_user = state.write().await.insert(new_id, ChatUser { name, password }); println!("{:?}", state); } }
方式二:每个连接只插入一次用户(业务一个连接对应一个用户)
用标志位控制只插入一次,避免重复key:
async fn web_socket(mut ws: WebSocket, State(state): State<UserDB>) { let new_id = NEW_USER_ID.fetch_add(1, std::sync::atomic::Ordering::Relaxed); let mut has_inserted = false; while let Ok(Message::Text(message)) = ws.next().await.unwrap(){ let user_info: Value = serde_json::from_str(&message.trim()).unwrap(); let (name, password) = (user_info["name"].to_string(), user_info["password"].to_string()); if !has_inserted { state.write().await.insert(new_id, ChatUser { name, password }); has_inserted = true; println!("{:?}", state); } // 后续消息可以做其他处理 } }
方式三:处理重复key场景(比如更新用户信息)
如果业务允许更新已存在的用户,或者需要明确处理重复情况,可以用HashMap的entryAPI:
async fn web_socket(mut ws: WebSocket, State(state): State<UserDB>) { let new_id = NEW_USER_ID.fetch_add(1, std::sync::atomic::Ordering::Relaxed); while let Ok(Message::Text(message)) = ws.next().await.unwrap(){ let user_info: Value = serde_json::from_str(&message.trim()).unwrap(); let (name, password) = (user_info["name"].to_string(), user_info["password"].to_string()); let mut db = state.write().await; // 使用entry API,存在则更新,不存在则插入 db.entry(new_id) .and_modify(|user| { user.name = name.clone(); user.password = password.clone(); }) .or_insert(ChatUser { name, password }); println!("{:?}", state); } }
额外建议
- 代码里的
unwrap()(包括ws.next().await.unwrap()和serde_json::from_str(...).unwrap())都可能触发panic,建议用match或if let处理错误,避免程序意外崩溃。 - 持有
RwLock写锁的时间尽量短,比如先解析完消息再获取锁,避免阻塞其他任务。
内容的提问来源于stack exchange,提问作者kecakisa
相关产品推荐
相关产品推荐

