如何创建贯穿程序生命周期的线程并传递不可变数据块处理?
问题:实时约束下线程初始化与工作负载拆分的生命周期问题
我有一批带实时约束的数学计算任务,主循环会反复调用相关函数并将结果存入现有缓冲区。希望在初始化阶段生成线程,线程处理任务后等待新数据,已用Barrier实现同步,但无法拆分线程生成与工作负载,尝试多种Arc或crossbeam用法均失败。
当前实现代码
pub const WORK_SIZE: usize = 524_288; pub const NUM_THREADS: usize = 6; pub const NUM_TASKS_PER_THREAD: usize = WORK_SIZE / NUM_THREADS; fn main() { let mut work: Vec<f64> = Vec::with_capacity(WORK_SIZE); for i in 0..WORK_SIZE { work.push(i as f64); } crossbeam::scope(|scope| { let threads: Vec<_> = work .chunks(NUM_TASKS_PER_THREAD) .map(|chunk| scope.spawn(move |_| chunk.iter().cloned().sum::<f64>())) .collect(); let threaded_time = std::time::Instant::now(); let thread_sum: f64 = threads.into_iter().map(|t| t.join().unwrap()).sum(); let threaded_micros = threaded_time.elapsed().as_micros() as f64; println!("threaded took: {:#?}", threaded_micros); let serial_time = std::time::Instant::now(); let no_thread_sum: f64 = work.iter().cloned().sum(); let serial_micros = serial_time.elapsed().as_micros() as f64; println!("serial took: {:#?}", serial_micros); assert_eq!(thread_sum, no_thread_sum); println!( "Threaded performace was {:?}", serial_micros / threaded_micros ); }) .unwrap(); }
期望实现的代码(存在生命周期问题)
use std::sync::{Arc, Barrier, Mutex}; use std::{slice::Chunks, thread::JoinHandle}; pub const WORK_SIZE: usize = 524_288; pub const NUM_THREADS: usize = 6; pub const NUM_TASKS_PER_THREAD: usize = WORK_SIZE / NUM_THREADS; //simplified version of what actual work that code base will do fn do_work(data: &[f64], result: Arc<Mutex<f64>>, barrier: Arc<Barrier>) { loop { barrier.wait(); let sum = data.into_iter().cloned().sum::<f64>(); let mut result = *result.lock().unwrap(); result += sum; } } fn init( mut data: Chunks<'_, f64>, result: &Arc<Mutex<f64>>, barrier: &Arc<Barrier>, ) -> Vec<std::thread::JoinHandle<()>> { let mut handles = Vec::with_capacity(NUM_THREADS); //spawn threads, in actual code these would be stored in a lib crate struct for i in 0..NUM_THREADS { let result = result.clone(); let barrier = barrier.clone(); let chunk = data.nth(i).unwrap(); handles.push(std::thread::spawn(|| { //Pass the particular thread the particular chunk it will operate on. do_work(chunk, result, barrier); })); } handles } fn main() { let mut work: Vec<f64> = Vec::with_capacity(WORK_SIZE); let mut result = Arc::new(Mutex::new(0.0)); for i in 0..WORK_SIZE { work.push(i as f64); } let work_barrier = Arc::new(Barrier::new(NUM_THREADS + 1)); let threads = init(work.chunks(NUM_THREADS), &result, &work_barrier); loop { work_barrier.wait(); //actual code base would do something with summation stored in result. println!("{:?}", result.lock().unwrap()); } }
核心问题
当前实现的核心问题是chunk的生命周期不足:Chunks产生的切片是借用自原Vec<f64>的,而线程会持有这个切片的引用超过原Vec的作用域(编译器无法保证原Vec在线程运行期间一直有效)。尝试用Arc包裹切片时,因为Arc需要持有拥有权,而切片是借用类型,直接Arc::new(chunk)无法满足生命周期要求。
解决方案
解决思路是让线程持有数据的拥有权而非借用,同时修正结果累加的逻辑错误。以下是修改后的完整代码:
use std::sync::{Arc, Barrier, Mutex}; use std::thread::JoinHandle; pub const WORK_SIZE: usize = 524_288; pub const NUM_THREADS: usize = 6; pub const NUM_TASKS_PER_THREAD: usize = WORK_SIZE / NUM_THREADS; // 线程执行的工作逻辑,持有数据的Arc拥有权 fn do_work(data: Arc<Vec<f64>>, result: Arc<Mutex<f64>>, barrier: Arc<Barrier>) { loop { barrier.wait(); let sum = data.iter().sum::<f64>(); // 直接修改锁内的值,而非副本 *result.lock().unwrap() += sum; } } fn init( data_chunks: Vec<Arc<Vec<f64>>>, result: &Arc<Mutex<f64>>, barrier: &Arc<Barrier>, ) -> Vec<JoinHandle<()>> { let mut handles = Vec::with_capacity(NUM_THREADS); for chunk in data_chunks { let result_clone = result.clone(); let barrier_clone = barrier.clone(); handles.push(std::thread::spawn(move || { do_work(chunk, result_clone, barrier_clone); })); } handles } fn main() { // 初始化数据并拆分为多个独立的子Vec,用Arc共享 let work: Vec<f64> = (0..WORK_SIZE).map(|i| i as f64).collect(); let data_chunks: Vec<Arc<Vec<f64>>> = work .chunks(NUM_TASKS_PER_THREAD) .map(|chunk| Arc::new(chunk.to_vec())) .collect(); let result = Arc::new(Mutex::new(0.0)); let work_barrier = Arc::new(Barrier::new(NUM_THREADS + 1)); // 初始化线程 let _threads = init(data_chunks, &result, &work_barrier); // 主循环:触发计算并处理结果 loop { // 重置结果(根据实际需求调整) *result.lock().unwrap() = 0.0; // 等待所有线程完成一轮计算 work_barrier.wait(); let final_sum = *result.lock().unwrap(); println!("当前计算结果: {}", final_sum); // 模拟实时任务的间隔(实际场景根据需求调整) std::thread::sleep(std::time::Duration::from_millis(100)); } }
关键修改点
- 生命周期问题解决:将原
Vec的切片转换为独立的Vec并包裹在Arc中,让线程持有Arc<Vec<f64>>的拥有权,彻底规避借用生命周期限制。 - 结果累加逻辑修正:直接修改
*result.lock().unwrap(),而非修改副本,确保累加操作生效。 - 主循环优化:增加结果重置逻辑,避免每次累加导致结果无限增大(可根据实际需求调整)。
如果原数据需要在主循环中动态更新,可以将Arc<Mutex<Vec<f64>>>作为共享数据结构,线程每次计算前从Mutex中获取最新切片,但这种方式会引入锁开销,需根据实时约束权衡。
内容的提问来源于stack exchange,提问作者Matthew Pittenger
相关产品推荐
相关产品推荐

