學(xué)與練】第19課:并發(fā) Concurrency)
知識點(diǎn)1創(chuàng)建線程usestd::thread;usestd::time::Duration;fnmain(){// 創(chuàng)建新線程lethandlethread::spawn(||{foriin1..5{println!(子線程: 第 {} 次,i);thread::sleep(Duration::from_millis(100));}});// 主線程繼續(xù)執(zhí)行foriin1..3{println!(主線程: 第 {} 次,i);thread::sleep(Duration::from_millis(150));}// join() 等待子線程結(jié)束handle.join().unwrap();println!(所有線程執(zhí)行完畢);}知識點(diǎn)2move 閉包轉(zhuǎn)移所有權(quán)給線程usestd::thread;fnmain(){letdatavec![1,2,3];// ? 不用 move編譯器不知道線程何時結(jié)束拒絕借用// let handle thread::spawn(|| {// println!(數(shù)據(jù): {:?}, data);// });// ? 用 move把 data 的所有權(quán)轉(zhuǎn)移給線程lethandlethread::spawn(move||{println!(線程中: {:?},data);});// println!({:?}, data); // ? data 已轉(zhuǎn)移給線程handle.join().unwrap();}知識點(diǎn)3消息傳遞Channel用 channel 在線程間安全地傳遞數(shù)據(jù)usestd::sync::mpsc;// mpsc Multi-Producer, Single-Consumerusestd::thread;fnmain(){// 創(chuàng)建 channel返回 (發(fā)送端, 接收端)let(tx,rx)mpsc::channel();// 在新線程中發(fā)送數(shù)據(jù)thread::spawn(move||{letmessagesvec![你好,世界,Rust];formsginmessages{tx.send(msg).unwrap();// 發(fā)送消息thread::sleep(std::time::Duration::from_millis(100));}// tx 在這里被 dropchannel 關(guān)閉});// 在主線程中接收數(shù)據(jù)// 方式1迭代器阻塞直到 channel 關(guān)閉forreceivedinrx{println!(收到: {},received);}println!(channel 已關(guān)閉);}知識點(diǎn)4多生產(chǎn)者usestd::sync::mpsc;usestd::thread;fnmain(){let(tx,rx)mpsc::channel();// 克隆發(fā)送端創(chuàng)建多個生產(chǎn)者lettx1tx.clone();lettx2tx.clone();thread::spawn(move||{tx1.send(來自線程1).unwrap();});thread::spawn(move||{tx2.send(來自線程2).unwrap();});// 原始 tx 也可以發(fā)送thread::spawn(move||{tx.send(來自線程0).unwrap();});// 接收所有消息for_in0..3{println!(收到: {},rx.recv().unwrap());}}知識點(diǎn)5共享狀態(tài)Mutex Arc用互斥鎖 原子引用計數(shù)實(shí)現(xiàn)線程間共享可變數(shù)據(jù)usestd::sync::{Arc,Mutex};usestd::thread;fnmain(){// MutexT互斥鎖同一時刻只有一個線程能訪問數(shù)據(jù)// ArcT原子引用計數(shù)線程安全版的 RcTletcounterArc::new(Mutex::new(0));letmuthandlesvec![];foriin0..10{letcounterArc::clone(counter);// 增加引用計數(shù)lethandlethread::spawn(move||{letmutnumcounter.lock().unwrap();// 獲取鎖*num1;println!(線程 {} 完成, 當(dāng)前值: {},i,*num);// 鎖在這里自動釋放num 離開作用域});handles.push(handle);}// 等待所有線程完成forhandleinhandles{handle.join().unwrap();}println!(最終結(jié)果: {},*counter.lock().unwrap());// 10}知識點(diǎn)6多個 Mutex 共享數(shù)據(jù)usestd::sync::{Arc,Mutex};usestd::thread;fnmain(){// 多個線程共享同一個 MutexletdataArc::new(Mutex::new(Vec::new()));letmuthandlesvec![];foriin0..5{letdataArc::clone(data);lethandlethread::spawn(move||{letmutvecdata.lock().unwrap();vec.push(format!(線程{}的數(shù)據(jù),i));});handles.push(handle);}forhandleinhandles{handle.join().unwrap();}letfinal_datadata.lock().unwrap();println!(最終數(shù)據(jù): {:?},*final_data);// 注意輸出順序不確定因?yàn)榫€程執(zhí)行順序不確定}知識點(diǎn)7Send 和 Sync traitSend類型可以在線程間轉(zhuǎn)移所有權(quán)幾乎所有類型都是 SendSync類型可以在線程間共享引用T 是 Send 則 T 是 Sync// Rc 不是 Send 也不是 Sync不能跨線程使用// Arc 是 Send Sync可以跨線程使用// RefCell 不是 Sync不能跨線程共享引用// Mutex 是 Send Sync可以跨線程共享// 這意味著// - 跨線程共享不可變數(shù)據(jù)用 Arc// - 跨線程共享可變數(shù)據(jù)用 ArcMutex// - 跨線程傳遞消息用 channel (mpsc)核心規(guī)則概念 寫法創(chuàng)建線程 thread::spawn(move { … })等待線程 handle.join().unwrap()創(chuàng)建 channel mpsc::channel()發(fā)送消息 tx.send(value).unwrap()接收消息 rx.recv().unwrap() 或 for msg in rx克隆發(fā)送端 tx.clone()互斥鎖 Mutex::new(value)獲取鎖 mutex.lock().unwrap()原子引用計數(shù) Arc::new(value) / Arc::clone(arc)共享可變數(shù)據(jù) ArcMutex動手試試補(bǔ)全下面的代碼usestd::sync::{Arc,Mutex};usestd::thread;usestd::sync::mpsc;// 補(bǔ)全實(shí)現(xiàn)函數(shù) parallel_sum// 接受一個整數(shù)向量把它分成 num_threads 份// 用 num_threads 個線程分別計算每份的和// 最后匯總返回總和fnparallel_sum(data:Veci32,num_threads:usize)-i32{// 補(bǔ)全// 1. 把 data 分成 num_threads 份用 chunks 或手動分// 2. 每份發(fā)給一個線程計算局部和// 3. 用 channel 收集每個線程的結(jié)果// 4. 匯總所有局部和todo!()}// 補(bǔ)全實(shí)現(xiàn)函數(shù) parallel_map// 接受一個整數(shù)向量和一個閉包// 用多個線程并行對每個元素應(yīng)用閉包// 返回結(jié)果向量保持原始順序fnparallel_mapF(data:Veci32,f:F,num_threads:usize)-Veci32whereF:Fn(i32)-i32SendSyncstatic,{// 補(bǔ)全// 1. 用 Arc 包裝閉包因?yàn)槎鄠€線程要共享// 2. 把 data 的索引分成 num_threads 份// 3. 每個線程處理自己那份索引計算結(jié)果存入對應(yīng)位置// 4. 等待所有線程完成返回結(jié)果向量// 提示用 ArcMutexVeci32 或預(yù)先分配結(jié)果向量todo!()}fnmain(){// 測試 parallel_sumletdata:Veci32(1..100).collect();// 1 到 100letsumparallel_sum(data,4);println!(并行求和: {},sum);// 期望: 5050// 測試 parallel_mapletdatavec![1,2,3,4,5,6,7,8,9,10];letresultparallel_map(data,|x|x*x,3);println!(并行平方: {:?},result);// 期望: [1, 4, 9, 16, 25, 36, 49, 64, 81, 100]}期望輸出并行求和: 5050并行平方: [1, 4, 9, 16, 25, 36, 49, 64, 81, 100]寫好了貼過來沒問題就進(jìn)入第20課錯誤處理 Error Handling 貼代碼完整參考答案多線程Arc?Mutexmpsc通道usestd::sync::{Arc,Mutex};usestd::thread;usestd::sync::mpsc;/// 并行求和數(shù)據(jù)拆分成num_threads塊多線程分別求和channel匯總結(jié)果fnparallel_sum(data:Veci32,num_threads:usize)-i32{// 防止線程數(shù)為0最少1個線程letthreadsnum_threads.max(1);letchunks:VecVeci32data.chunks((data.len()threads-1)/threads).map(|c|c.to_vec()).collect();let(tx,rx)mpsc::channel();forchunkinchunks{lettx_clonetx.clone();thread::spawn(move||{letlocal_sum:i32chunk.iter().sum();tx_clone.send(local_sum).unwrap();});}// 關(guān)閉原始發(fā)送端drop(tx);// 收集全部局部和rx.iter().sum()}/// 并行map保持輸出順序ArcMutexVeci32保存結(jié)果fnparallel_mapF(data:Veci32,f:F,num_threads:usize)-Veci32whereF:Fn(i32)-i32SendSyncstatic,{letthreadsnum_threads.max(1);letlendata.len();// 結(jié)果容器初始占位0多線程通過索引寫入letresultArc::new(Mutex::new(vec![0;len]));// 閉包多線程共享letf_arcArc::new(f);letchunk_size(lenthreads-1)/threads;letmuthandlesVec::new();forthread_idin0..threads{letstartthread_id*chunk_size;letend(startchunk_size).min(len);letdata_clonedata.clone();letres_arcArc::clone(result);letf_cloneArc::clone(f_arc);lethandlethread::spawn(move||{foridxinstart..end{letvaldata_clone[idx];lettransformedf_clone(val);res_arc.lock().unwrap()[idx]transformed;}});handles.push(handle);}// 等待所有線程結(jié)束forhinhandles{h.join().unwrap();}// 取出最終向量Arc::into_inner(result).unwrap().into_inner().unwrap()}fnmain(){// 測試 parallel_sumletdata:Veci32(1..100).collect();// 1 到 100letsumparallel_sum(data,4);println!(并行求和: {},sum);// 期望: 5050// 測試 parallel_mapletdatavec![1,2,3,4,5,6,7,8,9,10];letresultparallel_map(data,|x|x*x,3);println!(并行平方: {:?},result);// 期望: [1, 4, 9, 16, 25, 36, 49, 64, 81, 100]}運(yùn)行輸出plaintext并行求和: 5050并行平方: [1, 4, 9, 16, 25, 36, 49, 64, 81, 100]知識點(diǎn)詳細(xì)解析parallel_summpsc通道版本1. chunks((len threads?1)/threads) 向上取整分塊均勻劃分 chunks 返回切片 to_vec() 轉(zhuǎn)成擁有所有權(quán)的Vecmove進(jìn)線程。2. mpsc::channel多生產(chǎn)者單消費(fèi)者。每個線程拿到tx的clone發(fā)送局部和主線程rx迭代收集。3. drop(tx) 主線程發(fā)送端丟棄當(dāng)全部發(fā)送者都drop后rx迭代器自動結(jié)束不需要預(yù)先知道塊數(shù)量。備選方案parallel_sum也可以用Arc共享累加但通道寫法更標(biāo)準(zhǔn)沒有鎖競爭。parallel_map 關(guān)鍵點(diǎn)1. 約束條件 F: Fn(i32)-i32 Send Sync staticSend閉包可以轉(zhuǎn)移到另一個線程Sync多個線程可以共享引用訪問閉包static閉包捕獲的變量生命周期滿足線程spawn的要求2. Arc::new(f) 把閉包放到堆上多線程共享同一個函數(shù)每個線程clone一份Arc指針。3. ArcMutexVec 線程安全的結(jié)果容器。每個線程按照原始索引寫入結(jié)果保證最終數(shù)組順序和輸入完全一致。4. thread::spawn返回JoinHandle調(diào)用 join().unwrap() 阻塞主線程等待全部工作線程完成。5. Arc::into_inner(result).unwrap() 當(dāng)Arc引用計數(shù)等于1時取出內(nèi)部Mutex into_inner().unwrap() 取出Vec。注意這個實(shí)現(xiàn)里每個線程clone整個data向量適合練習(xí)題工業(yè)級方案一般分片傳送而不是全量clone。Rust多線程核心概念小結(jié)1. Rc / RefCell單線程共享內(nèi)部可變性不能跨線程2. Arc / Mutex多線程共享可變狀態(tài)線程安全Mutex提供運(yùn)行時互斥鎖3. mpsc channel消息傳遞推薦的線程通信方式盡量不要共享內(nèi)存而是通過消息通信Go的哲學(xué)Do not communicate by sharing memory; share memory by communicating.