Rust多线程共享Vec数据更新失效问题排查与修复
Rust多线程共享Vec<[u8;128]>更新无效的修复方案
问题根源
- 更新逻辑完全错误:原
update方法仅修改了Vec元素的副本,未触及原数据。代码中通过读锁取出元素后,创建了一个新的RwLock修改副本,这种操作对共享Vec没有任何影响。 - 读写锁死锁风险:
process函数先持有读锁,再调用update尝试获取写锁,会触发读写锁的死锁(读锁未释放时无法获取写锁)。 - 结构冗余:
SharedData内部的data已被Arc<RwLock>包裹,外层又嵌套Arc<SharedData>,虽不影响功能,但可简化以提高可读性。
修复后的完整代码
main.rs
use std::net::SocketAddr; use std::sync::Arc; use std::sync::RwLock; use std::time::Duration; use tokio_task_pool::Pool; struct SharedData { data: Arc<RwLock<Vec<[u8; 128]>>> } impl SharedData { // 直接接收初始Vec,内部封装为Arc<RwLock> fn new(initial_data: Vec<[u8; 128]>) -> Self { Self { data: Arc::new(RwLock::new(initial_data)) } } // 正确的更新逻辑:获取写锁直接修改原Vec元素 fn update(&self, index: usize, update_data: [u8; 128]) { let mut write_guard = self.data.write().unwrap(); // 增加索引越界判断,避免panic if index < write_guard.len() { write_guard[index] = update_data; } } // 封装读取逻辑,自动管理锁的生命周期 fn get_current_data(&self) -> Vec<[u8; 128]> { let read_guard = self.data.read().unwrap(); read_guard.clone() } } fn socket_to_async_tcplistener(s: socket2::Socket) -> std::io::Result<tokio::net::TcpListener> { std::net::TcpListener::from(s).try_into() } async fn process(mut stream: tokio::net::TcpStream, db_arc: Arc<SharedData>) { // 先读取数据,锁在clone后自动释放 let current_data = db_arc.get_current_data(); println!("In process() read: {:?}", current_data); // 读取完成后再执行更新,避免死锁 db_arc.update(1, [3u8; 128]); } async fn serve(_: usize, tcplistener_arc: Arc<tokio::net::TcpListener>, db_arc: Arc<SharedData>) { let task_pool_capacity = 10; let task_pool = Pool::bounded(task_pool_capacity) .with_spawn_timeout(Duration::from_secs(300)) .with_run_timeout(Duration::from_secs(300)); loop { let (stream, _) = tcplistener_arc.as_ref().accept().await.unwrap(); let db_arc_clone = db_arc.clone(); task_pool.spawn(async move { process(stream, db_arc_clone).await; }).await.unwrap(); } } #[tokio::main] async fn main() { let addr: std::net::SocketAddr = "0.0.0.0:50051".parse().unwrap(); let soc2 = socket2::Socket::new( match addr { SocketAddr::V4(_) => socket2::Domain::IPV4, SocketAddr::V6(_) => socket2::Domain::IPV6, }, socket2::Type::STREAM, Some(socket2::Protocol::TCP) ).unwrap(); soc2.set_reuse_address(true).unwrap(); soc2.set_reuse_port(true).unwrap(); soc2.set_nonblocking(true).unwrap(); soc2.bind(&addr.into()).unwrap(); soc2.listen(8192).unwrap(); let tcp_listener = Arc::new(socket_to_async_tcplistener(soc2).unwrap()); let initial_vec = vec![ [0u8; 128], [1u8; 128], [2u8; 128], ]; let share_db = Arc::new(SharedData::new(initial_vec)); let mut handlers = Vec::new(); for i in 0..num_cpus::get() - 1 { let cloned_listener = Arc::clone(&tcp_listener); let db_arc = share_db.clone(); let h = std::thread::spawn(move || { tokio::runtime::Builder::new_current_thread() .enable_all() .build() .unwrap() .block_on(serve(i, cloned_listener, db_arc)); }); handlers.push(h); } for h in handlers { h.join().unwrap(); } }
Cargo.toml(无需修改)
[package] name = "tokio-test" version = "0.1.0" edition = "2021" [dependencies] log = "0.4.20" env_logger = "0.10.0" tokio = { version = "1.34.0", features = ["full"] } tokio-stream = { version = "0.1.14", features = ["net"] } serde = { version = "1.0.193", features = ["derive"] } serde_yaml = "0.9.27" serde_derive = "1.0.193" mio = {version="0.8.9", features=["net", "os-poll", "os-ext"]} num_cpus = "1.16.0" socket2 = { version="0.5.5", features = ["all"]} array-macro = "2.1.8" tokio-task-pool = "0.1.5" argparse = "0.2.2"
关键修改说明
- 修正update方法:通过获取
RwLock的写权限,直接修改共享Vec中的目标元素,确保修改作用于原数据。 - 避免死锁:将读取和更新操作分离,读取完成后释放读锁,再执行更新操作,消除读写锁冲突。
- 封装读写逻辑:新增
get_current_data方法封装读取操作,自动管理锁的生命周期,减少手动操作锁的出错概率。 - 简化初始化流程:
SharedData::new直接接收初始Vec,内部完成Arc<RwLock>的封装,代码更简洁。
内容的提问来源于stack exchange,提问作者matsuisa
相关产品推荐
相关产品推荐

