在Rust中动态创建多Mutex实现主从线程数据同步的方案咨询
运行时动态创建线程并共享矩阵数据的Rust实现
一、基于Arc的共享方案
你可以通过创建一个包含n个Arc<Mutex<Option<Matrix>>>的向量来实现动态共享,每个元素对应一个从线程的最新矩阵数据。Arc用于跨线程共享所有权,Mutex保证线程安全的访问,Option则处理“尚未输入数据”的初始状态。
示例代码
use std::sync::{Arc, Mutex}; use std::thread; use std::io::{self, BufRead}; // 定义矩阵类型,简化后续代码 type Matrix = Vec<Vec<i32>>; fn main() { // 运行时获取线程数量n println!("请输入线程数量n:"); let mut input = String::new(); io::stdin().lock().read_line(&mut input).unwrap(); let n: usize = input.trim().parse().unwrap(); // 创建存储每个线程最新矩阵的共享向量 let shared_matrices: Vec<Arc<Mutex<Option<Matrix>>>> = (0..n) .map(|_| Arc::new(Mutex::new(None))) .collect(); // 启动n个从线程 let mut handles = vec![]; for i in 0..n { let matrix_ptr = Arc::clone(&shared_matrices[i]); let handle = thread::spawn(move || loop { println!("\n线程{}: 请输入矩阵(每行数字用空格分隔,空行结束输入):", i); let mut lines = vec![]; let stdin = io::stdin().lock(); for line in stdin.lines() { let line = line.unwrap(); let trimmed = line.trim(); if trimmed.is_empty() { break; } // 解析行数据为整数向量 let row: Vec<i32> = trimmed .split_whitespace() .map(|s| s.parse().unwrap()) .collect(); lines.push(row); } if !lines.is_empty() { // 更新共享的最新矩阵 let mut matrix = matrix_ptr.lock().unwrap(); *matrix = Some(lines); println!("线程{}: 矩阵已更新", i); } }); handles.push(handle); } // 主线程持续汇总并处理最新数据 loop { println!("\n=== 主线程汇总最新矩阵 ==="); for (i, mat_ptr) in shared_matrices.iter().enumerate() { let matrix = mat_ptr.lock().unwrap(); match &*matrix { Some(mat) => { println!("线程{}的最新矩阵:", i); for row in mat { println!("{:?}", row); } // 这里可以添加你的额外操作,比如计算矩阵的和、转置等 let sum: i32 = mat.iter().flat_map(|r| r.iter()).sum(); println!("线程{}矩阵元素总和: {}", i, sum); } None => println!("线程{}尚未输入矩阵", i), } } // 暂停1秒,避免频繁输出 thread::sleep(std::time::Duration::from_secs(1)); } // 注:实际场景中需要处理线程退出逻辑,这里为了简化省略了join // for handle in handles { // handle.join().unwrap(); // } }
关键说明
shared_matrices是动态创建的向量,每个元素都是Arc<Mutex<Option<Matrix>>>,对应一个线程的共享数据- 每个从线程克隆对应的
Arc,保证所有权跨线程共享 Mutex::lock()会阻塞直到获取锁,确保同一时间只有一个线程修改或读取数据
二、替代方案:基于消息通道(mpsc)
如果想避免显式使用锁,可以用Rust标准库的mpsc消息通道。每个从线程将最新矩阵发送给主线程,主线程维护一个数组存储每个线程的最新数据,这种方式更符合Rust的“所有权转移”模型,减少竞态风险。
示例代码
use std::sync::mpsc; use std::thread; use std::io::{self, BufRead}; type Matrix = Vec<Vec<i32>>; fn main() { println!("请输入线程数量n:"); let mut input = String::new(); io::stdin().lock().read_line(&mut input).unwrap(); let n: usize = input.trim().parse().unwrap(); // 创建n个通道,每个线程对应一个sender,主线程持有所有receiver let mut senders = vec![]; let mut receivers = vec![]; for _ in 0..n { let (tx, rx) = mpsc::channel(); senders.push(tx); receivers.push(rx); } // 启动n个从线程 let mut handles = vec![]; for i in 0..n { let tx = senders[i].clone(); let handle = thread::spawn(move || loop { println!("\n线程{}: 请输入矩阵(每行数字用空格分隔,空行结束输入):", i); let mut lines = vec![]; let stdin = io::stdin().lock(); for line in stdin.lines() { let line = line.unwrap(); let trimmed = line.trim(); if trimmed.is_empty() { break; } let row: Vec<i32> = trimmed .split_whitespace() .map(|s| s.parse().unwrap()) .collect(); lines.push(row); } if !lines.is_empty() { // 发送最新矩阵到主线程 tx.send((i, lines)).unwrap(); println!("线程{}: 矩阵已发送", i); } }); handles.push(handle); } // 主线程维护每个线程的最新矩阵 let mut latest_matrices: Vec<Option<Matrix>> = vec![None; n]; println!("\n=== 主线程开始监听矩阵更新 ==="); // 使用select同时监听所有receiver(标准库mpsc不支持直接select,这里用循环轮询) loop { for rx in receivers.iter_mut() { if let Ok((thread_id, matrix)) = rx.try_recv() { latest_matrices[thread_id] = Some(matrix); println!("\n线程{}的矩阵已更新,当前汇总数据:", thread_id); // 展示所有最新矩阵并执行额外操作 for (i, mat) in latest_matrices.iter().enumerate() { match mat { Some(m) => { println!("线程{}的矩阵:", i); for row in m { println!("{:?}", row); } let sum: i32 = m.iter().flat_map(|r| r.iter()).sum(); println!("元素总和: {}", sum); } None => println!("线程{}无数据", i), } } } } thread::sleep(std::time::Duration::from_millis(500)); } // 省略join逻辑 // for handle in handles { // handle.join().unwrap(); // } }
关键说明
- 每个线程对应一个独立的通道,避免消息混淆
- 主线程通过
try_recv()非阻塞轮询所有通道,更新最新数据 - 无需手动处理锁,通道天然保证线程安全
内容的提问来源于stack exchange,提问作者Dev. fennek
相关产品推荐
相关产品推荐

