Channels in Tokio
Tokio provides async channels for task coordination: mpsc for work queues, oneshot for single replies, broadcast for fan-out, and watch for latest-value signaling.
Busca en todas las páginas de la documentación
Tokio provides async channels for task coordination: mpsc for work queues, oneshot for single replies, broadcast for fan-out, and watch for latest-value signaling.
Quick-reference recipe card - copy-paste ready.
use tokio::sync::mpsc;
#[tokio::main]
async fn main() {
let (tx, mut rx) = mpsc::channel(32);
tx.send("job").await.unwrap();
println!("{}", rx.recv().await.unwrap());
}When to reach for this: Worker pools, request/response between tasks, shutdown broadcasts, and config change notifications.
use tokio::sync::{broadcast, oneshot, watch};
#[tokio::main]
async fn main() {
// oneshot: single response
let (otx, orx) = oneshot::channel();
tokio::spawn(async move {
otx.send(42).unwrap();
});
println!("oneshot: {}", orx.await.unwrap());
// watch: latest value
let (wtx, mut wrx) = watch::channel("v1".to_string());
wtx.send("v2".into()).unwrap();
wrx.changed().await.unwrap();
println!("watch: {}", *wrx.borrow());
// broadcast: all subscribers
let (btx, mut brx) = broadcast::channel(16);
btx.send("event").unwrap();
println!("broadcast: {}", brx.recv().await.unwrap());
}What this demonstrates:
watch overwrites - only latest value matters.broadcast may lag if subscribers are slow (configurable error).send().await applies backpressure.Sender::send updates; Receiver::changed().await waits for new value.| Channel | Messages | Consumers | Typical use |
|---|---|---|---|
| mpsc | many | one | job queue |
| oneshot | one | one | single reply |
| broadcast | many | many | events, ticks |
| watch | latest | many | config, shutdown flag |
use tokio_stream::wrappers::ReceiverStream;
use tokio_stream::StreamExt;
let (tx, rx) = mpsc::channel(8);
let mut stream = ReceiverStream::new(rx);
while let Some(item) = stream.next().await {
println!("{item}");
}mpsc::Receiver to Stream for iterator-style processing.Semaphore complements channels for limiting in-flight work.mpsc over Mutex<VecDeque> for task boundaries.channel(1) with slow consumer still blocks at capacity; unbounded not in std Tokio mpsc (always bounded). Fix: Size buffer to expected burst.send fails. Fix: One receiver only; clone data before send if needed elsewhere.RecvError::Lagged.changed().await before relying on updates in loops.
| Alternative | Use When | Don't Use When |
|---|---|---|
std::sync::mpsc | Sync threads only | Async tasks awaiting recv |
crossbeam | Sync MPMC patterns | Pure async await loops |
Mutex + Condvar | Legacy sync | Idiomatic Tokio services |
| Actor crate | Complex state machines | Simple request/response |
Match burst size and backpressure tolerance - too small stalls producers; too large hides overload.
oneshot is lighter for single reply; mpsc when multiple messages possible.
broadcast delivers every event; watch keeps only latest snapshot.
Drop all Sender handles; recv returns None.
Non-async poll - use in select! or bridging sync code carefully.
mpsc FIFO per sender interleaving - not strict global fairness across many producers.
Spawn worker task with mpsc receiver; handlers send jobs via cloned Sender in state.
mpsc::recv in select! loops is cancel-safe per Tokio docs - still design for lost wakeups in complex setups.
Semaphore limits concurrency; channel also transports data - often used together.
Use small buffers to force backpressure in tests; timeout on recv to catch hangs.
Stack versions: This page was written for Rust 1.97.0 (edition 2024), Tokio 1.x, Axum 0.8, serde 1.0, sqlx 0.8, clap 4, and Polars 0.46+.
Revisado por Chris St. John·Última actualización: 19 jul 2026