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
Arc<Mutex<i64>>shared, lockedthread::spawn(move Rc<T>not Send: rejectedSendsafe to move to another threadSyncsafe to share references between threadschannelsmpsc, crossbeam, tokio::sync
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
async fn handler()compiles to a state machineFuturedoes nothing until polledTokio runtimemulti-threaded work-stealingwakers"poll me again"reactorepoll / kqueue / io_uringdrop = cancela dropped future stops
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