Rust 中的 Tokio 線程同步機(jī)制詳解
Rust 中的 Tokio 線程同步機(jī)制
在并發(fā)編程中,線程同步是一個(gè)重要的概念,用于確保多個(gè)線程在訪問(wèn)共享資源時(shí)能夠正確地協(xié)調(diào)。Tokio 是一個(gè)強(qiáng)大的異步運(yùn)行時(shí)庫(kù),為 Rust 提供了多種線程同步機(jī)制。以下是一些常見的同步機(jī)制:
- Mutex
- RwLock
- Barrier
- Semaphore
- Notify
- oneshot 和 mpsc 通道
- watch 通道
1. Mutex
Mutex(互斥鎖)是最常見的同步原語(yǔ)之一,用于保護(hù)共享數(shù)據(jù)。它確保同一時(shí)間只有一個(gè)線程能夠訪問(wèn)數(shù)據(jù),從而避免競(jìng)爭(zhēng)條件。
use tokio::sync::Mutex;
use std::sync::Arc;
?
#[tokio::main]
async fn main() {
let data = Arc::new(Mutex::new(0));
?
let mut handles = vec![];
for _ in 0..10 {
let data = data.clone();
let handle = tokio::spawn(async move {
let mut lock = data.lock().await;
*lock += 1;
});
handles.push(handle);
}
?
for handle in handles {
handle.await.unwrap();
}
?
println!("Result: {}", *data.lock().await);
}2. RwLock
RwLock(讀寫鎖)允許多線程同時(shí)讀取數(shù)據(jù),但只允許一個(gè)線程寫入數(shù)據(jù)。它比 Mutex 更加靈活,因?yàn)樵谧x取多于寫入的場(chǎng)景下,它能提高性能。功能上,他是讀寫互斥、寫寫互斥、讀讀兼容。
use tokio::sync::RwLock;
use std::sync::Arc;
?
#[tokio::main]
async fn main() {
let data = Arc::new(RwLock::new(0));
?
let read_data = data.clone();
let read_handle = tokio::spawn(async move {
let lock = read_data.read().await;
println!("Read: {}", *lock);
});
?
let write_data = data.clone();
let write_handle = tokio::spawn(async move {
let mut lock = write_data.write().await;
*lock += 1;
println!("Write: {}", *lock);
});
?
read_handle.await.unwrap();
write_handle.await.unwrap();
}3. Barrier
Barrier 是一種同步機(jī)制,允許多個(gè)線程在某個(gè)點(diǎn)上進(jìn)行同步。當(dāng)線程到達(dá)屏障時(shí),它們會(huì)等待直到所有線程都到達(dá),然后一起繼續(xù)執(zhí)行。
use tokio::sync::Barrier;
use std::sync::Arc;
?
#[tokio::main]
async fn main() {
let barrier = Arc::new(Barrier::new(3));
?
let mut handles = vec![];
for i in 0..3 {
let barrier = barrier.clone();
let handle = tokio::spawn(async move {
println!("Before wait: {}", i);
barrier.wait().await;
println!("After wait: {}", i);
});
handles.push(handle);
}
?
for handle in handles {
handle.await.unwrap();
}
}4. Semaphore
Semaphore(信號(hào)量)是一種用于控制對(duì)資源訪問(wèn)的同步原語(yǔ)。它允許多個(gè)線程訪問(wèn)資源,但有一個(gè)最大并發(fā)數(shù)限制。
#[tokio::test]
async fn test_sem() {
let semaphore = Arc::new(Semaphore::new(3));
?
let mut handles = vec![];
for i in 0..5 {
let semaphore = semaphore.clone();
let handle = tokio::spawn(async move {
let permit = semaphore.acquire().await.unwrap();
let now = Local::now();
println!("Got permit: {} at {:?}", i, now);
println!(
"Semaphore available permits before sleep: {}",
semaphore.available_permits()
);
sleep(Duration::from_secs(5)).await;
drop(permit);
println!(
"Semaphore available permits after sleep: {}",
semaphore.available_permits()
);
});
handles.push(handle);
}
?
for handle in handles {
handle.await.unwrap();
}
}最終的結(jié)果如下
Got permit: 0 at 2024-08-08T21:03:04.374666+08:00
Semaphore available permits before sleep: 2
Got permit: 1 at 2024-08-08T21:03:04.375527800+08:00
Semaphore available permits before sleep: 1
Got permit: 2 at 2024-08-08T21:03:04.375563+08:00
Semaphore available permits before sleep: 0
Semaphore available permits after sleep: 0
Semaphore available permits after sleep: 0
Semaphore available permits after sleep: 1
Got permit: 3 at 2024-08-08T21:03:09.376722800+08:00
Semaphore available permits before sleep: 1
Got permit: 4 at 2024-08-08T21:03:09.376779200+08:00
Semaphore available permits before sleep: 1
Semaphore available permits after sleep: 2
Semaphore available permits after sleep: 3
5. Notify
Notify 是一種用于線程間通知的簡(jiǎn)單機(jī)制。它允許一個(gè)線程通知其他線程某些事件的發(fā)生。
use tokio::sync::Notify;
use std::sync::Arc;
?
#[tokio::main]
async fn main() {
let notify = Arc::new(Notify::new());
let notify_clone = notify.clone();
?
let handle = tokio::spawn(async move {
notify_clone.notified().await;
println!("Received notification");
});
?
notify.notify_one();
handle.await.unwrap();
}6. oneshot 和 mpsc 通道
oneshot 通道用于一次性發(fā)送消息,而 mpsc 通道則允許多個(gè)生產(chǎn)者發(fā)送消息到一個(gè)消費(fèi)者。一般地onshot用于異常通知、啟動(dòng)分析等功能。mpsc用于實(shí)現(xiàn)異步消息同步
oneshot
use tokio::sync::oneshot;
?
#[tokio::main]
async fn main() {
let (tx, rx) = oneshot::channel();
?
tokio::spawn(async move {
tx.send("Hello, world!").unwrap();
});
?
let message = rx.await.unwrap();
println!("Received: {}", message);
}mpsc
use tokio::sync::mpsc;
?
#[tokio::main]
async fn main() {
let (tx, mut rx) = mpsc::channel(32);
?
tokio::spawn(async move {
tx.send("Hello, world!").await.unwrap();
});
?
while let Some(message) = rx.recv().await {
println!("Received: {}", message);
}
}7. watch 通道
watch 通道用于發(fā)送和接收共享狀態(tài)的更新。它允許多個(gè)消費(fèi)者監(jiān)聽狀態(tài)的變化。
use tokio::sync::watch;
?
#[tokio::main]
async fn main() {
let (tx, mut rx) = watch::channel("initial");
?
tokio::spawn(async move {
tx.send("updated").unwrap();
});
?
while rx.changed().await.is_ok() {
println!("Received: {}", *rx.borrow());
}
}?watch通道?:
- 用于廣播狀態(tài)更新,一個(gè)生產(chǎn)者更新狀態(tài),多個(gè)消費(fèi)者獲取最新狀態(tài)。
- 適合配置變更、狀態(tài)同步等場(chǎng)景。
?mpsc通道?:
- 用于傳遞消息隊(duì)列,多個(gè)生產(chǎn)者發(fā)送消息,一個(gè)消費(fèi)者逐條處理。
- 適合任務(wù)隊(duì)列、事件驅(qū)動(dòng)等場(chǎng)景。
總結(jié)
Rust 中的 Tokio 提供了豐富的線程同步機(jī)制,可以根據(jù)具體需求選擇合適的同步原語(yǔ)。常用的同步機(jī)制包括:
Mutex:互斥鎖,保護(hù)共享數(shù)據(jù)。RwLock:讀寫鎖,允許并發(fā)讀,寫時(shí)獨(dú)占。Barrier:屏障,同步多個(gè)線程在某一點(diǎn)。Semaphore:信號(hào)量,控制并發(fā)訪問(wèn)資源。Notify:通知機(jī)制,用于線程間通知。oneshot和mpsc通道:消息傳遞機(jī)制。watch通道:狀態(tài)更新機(jī)制。
通過(guò)這些同步機(jī)制,可以在 Rust 中編寫高效、安全的并發(fā)程序。
到此這篇關(guān)于Rust 中的 Tokio 線程同步機(jī)制的文章就介紹到這了,更多相關(guān)rust tokio線程同步內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Rust開發(fā)WebAssembly在Html和Vue中的應(yīng)用小結(jié)(推薦)
這篇文章主要介紹了Rust開發(fā)WebAssembly在Html和Vue中的應(yīng)用,本文將帶領(lǐng)大家在普通html上和vue手腳架上都來(lái)運(yùn)行wasm的流程,需要的朋友可以參考下2022-08-08
使用?Rust?實(shí)現(xiàn)的基礎(chǔ)的List?和?Watch?機(jī)制示例流程
本文給大家介紹使用Rust實(shí)現(xiàn)的基礎(chǔ)的List和Watch機(jī)制示例流程,本文通過(guò)實(shí)例代碼給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友參考下吧2025-09-09
Rust生命周期之驗(yàn)證引用有效性與防止懸垂引用方式
本文介紹了Rust中生命周期注解的應(yīng)用,包括防止懸垂引用、在函數(shù)中使用泛型生命周期、生命周期省略規(guī)則、在結(jié)構(gòu)體中使用生命周期、靜態(tài)生命周期以及如何將生命周期與泛型和特質(zhì)約束結(jié)合,通過(guò)這些機(jī)制,Rust在編譯時(shí)就能捕獲內(nèi)存安全問(wèn)題2025-02-02

