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

Rust多线程共享Vec数据更新失效问题排查与修复

Rust多线程共享Vec<[u8;128]>更新无效的修复方案

问题根源

  1. 更新逻辑完全错误:原update方法仅修改了Vec元素的副本,未触及原数据。代码中通过读锁取出元素后,创建了一个新的RwLock修改副本,这种操作对共享Vec没有任何影响。
  2. 读写锁死锁风险:process函数先持有读锁,再调用update尝试获取写锁,会触发读写锁的死锁(读锁未释放时无法获取写锁)。
  3. 结构冗余: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"

关键修改说明

  1. 修正update方法:通过获取RwLock的写权限,直接修改共享Vec中的目标元素,确保修改作用于原数据。
  2. 避免死锁:将读取和更新操作分离,读取完成后释放读锁,再执行更新操作,消除读写锁冲突。
  3. 封装读写逻辑:新增get_current_data方法封装读取操作,自动管理锁的生命周期,减少手动操作锁的出错概率。
  4. 简化初始化流程:SharedData::new直接接收初始Vec,内部完成Arc<RwLock>的封装,代码更简洁。

内容的提问来源于stack exchange,提问作者matsuisa

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 03:45:54