Part 4 · 2 chapters · ~12 min
Concurrency and Async
Threads, Send and Sync, Arc and Mutex with measured costs, channels, rayon for data parallelism, async/await as lazy futures, the Tokio runtime, wakers and the reactor, cancellation by drop and cancel safety, select and timeouts, spawn_blocking, and structured task management with JoinSet.
7
Threads, Send and Sync
code
// measured on this machine: spawn + join of 2,000 OS threads ≈ 35.7 µs each // 8 threads × 100,000 increments on Arc<Mutex<i64>> → exactly 800,000 in 17.0 ms use rayon::prelude::*; let total: i64 = transfers.par_iter().map(|t| t.amount_kobo).sum(); // data parallelism over all cores
FEARLESS CONCURRENCY
Send, Sync and the compiler refusing data races
swipe the figure sideways, or tap expand for full screen
1/5
the marker traits
Send means a value can be moved to another thread; Sync means references to it can be shared between threads. The compiler derives them automatically from a type's fields.
Send: move across threads; Sync: share across threadsauto-derived from fields
8
Tokio
code
#[tokio::main]
async fn main() -> anyhow::Result<()> {
let mut set = tokio::task::JoinSet::new();
for id in account_ids { set.spawn(async move { fetch_balance(id).await }); }
while let Some(res) = set.join_next().await { println!("{:?}", res??); }
match tokio::time::timeout(Duration::from_secs(2), rail.submit(&payout)).await {
Ok(Ok(status)) => record(status),
Ok(Err(e)) => fail(e),
Err(_elapsed) => mark_unknown(&payout), // the rail call was dropped (cancelled) at its await point
}
let pdf = tokio::task::spawn_blocking(move || render_statement(&stmt)).await?; // CPU work off the runtime
Ok(())
}ASYNC RUST WITH TOKIO
futures are lazy state machines; Tokio polls them
swipe the figure sideways, or tap expand for full screen
1/5
lazy futures
An async fn returns a Future that does nothing until it is polled. The compiler turns it into a state machine whose states are the await points (Concurrency part 5).
nothing runs until polleda compiler-generated state machine