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

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 04:44:55