Q: Channels in Rust — mpsc, oneshot, broadcast, watch. Which one when?

Answer:

Rust ecosystem ships several channel types. They look similar (send/recv) but have different semantics around senders, receivers, capacity, and replay. Picking the wrong one is the most common cause of "the message isn't being delivered" bugs.

The Cheat Sheet

ChannelSendersReceiversCapacitySemantics
mpscmanyonebounded or unboundedFIFO queue, latest-receiver wins
spsconeoneboundedFaster, single-producer
oneshotoneone1Send once, receive once
broadcastmanymanybounded ringEach receiver gets every message
watchonemany1 (latest only)New subscribers see latest
bounded (crossbeam)manymanyboundedmpmc

mpsc — Multi-Producer, Single-Consumer

Tokio's tokio::sync::mpsc is the workhorse for async pipelines.

#![allow(unused)]
fn main() {
use tokio::sync::mpsc;

let (tx, mut rx) = mpsc::channel::<Event>(100);   // bounded

// Many producers
for _ in 0..4 {
    let tx = tx.clone();
    tokio::spawn(async move {
        tx.send(Event::new()).await.unwrap();
    });
}

// One consumer
while let Some(event) = rx.recv().await {
    process(event).await;
}
}

Bounded channels apply backpressuresend() awaits until there's space. This is the whole point: a slow consumer slows producers, not OOMs the queue.

Unbounded variant:

#![allow(unused)]
fn main() {
let (tx, mut rx) = mpsc::unbounded_channel();
tx.send(event).unwrap();             // non-async send, no backpressure
}

Avoid unbounded for production. They convert "consumer is slow" into "process runs out of memory and crashes."

oneshot — Single-Use Reply

For request/response patterns where one task asks another for one value:

#![allow(unused)]
fn main() {
use tokio::sync::oneshot;

let (resp_tx, resp_rx) = oneshot::channel::<Reply>();
request_tx.send(Req { resp: resp_tx, ... }).await?;

let reply = resp_rx.await?;          // wait for one reply
}

The handler holds resp_tx and calls resp_tx.send(reply) exactly once. After that, both ends are consumed.

Use cases:

  • Per-request reply channels.
  • "Done" signals from spawned tasks.
  • One-shot cancellation.

broadcast — Pub/Sub

Every subscriber receives every message after they subscribe. Like a Linux ring buffer.

#![allow(unused)]
fn main() {
use tokio::sync::broadcast;

let (tx, _) = broadcast::channel::<Event>(32);

// Subscribers must call subscribe() and start receiving BEFORE messages send,
// or they'll miss them.
let mut rx1 = tx.subscribe();
let mut rx2 = tx.subscribe();

tokio::spawn(async move {
    tx.send(event).unwrap();
});

let _ = rx1.recv().await;
let _ = rx2.recv().await;
}

Important: capacity is per-subscriber. If a slow subscriber falls behind by more than the capacity, it gets a RecvError::Lagged(N) indicating it missed N messages — the channel doesn't block the producer for one slow consumer.

Use for: shutdown signals, config updates, event broadcasting.

watch — Latest-Value Subject

A single-slot value with version tracking. Receivers always see the latest value, not the history.

#![allow(unused)]
fn main() {
use tokio::sync::watch;

let (tx, mut rx) = watch::channel(Config::default());

// Producer
tx.send(new_config).unwrap();

// Consumers
let cfg = rx.borrow().clone();        // current value, no wait
rx.changed().await?;                  // wait until next change
}

Use for: configuration updates, leader election state, "current value" semantics. Late subscribers see the current value, not the history.

Crossbeam Channels

Outside Tokio (or for sync code):

#![allow(unused)]
fn main() {
use crossbeam_channel::{bounded, unbounded, select};

let (s, r) = bounded::<i32>(100);

// MPMC — many senders, many receivers
let r2 = r.clone();
let s2 = s.clone();
}

Crossbeam supports select! macro over multiple channels (Go-style):

#![allow(unused)]
fn main() {
select! {
    recv(r1) -> msg => handle(msg.unwrap()),
    recv(r2) -> msg => handle(msg.unwrap()),
    default(Duration::from_secs(1)) => println!("timeout"),
}
}

std::sync::mpsc (Avoid)

The standard library's mpsc exists for historical reasons. It's slower than crossbeam, has fewer features, and isn't async-aware. Use tokio::sync::mpsc (async) or crossbeam_channel (sync).

Backpressure Patterns

Bounded channels are the easy 80%. For more shape:

1. Drop-when-full (lose messages instead of blocking producer):

#![allow(unused)]
fn main() {
match tx.try_send(event) {
    Ok(_)                        => {},
    Err(TrySendError::Full(_))   => metrics.dropped.inc(),
    Err(TrySendError::Closed(_)) => break,
}
}

Used for telemetry where freshness > completeness.

2. Coalesce / batch (collect N before processing):

#![allow(unused)]
fn main() {
let mut batch = Vec::new();
loop {
    tokio::select! {
        Some(item) = rx.recv() => {
            batch.push(item);
            if batch.len() >= 100 { flush(&mut batch).await; }
        }
        _ = tokio::time::sleep(Duration::from_secs(1)) => {
            if !batch.is_empty() { flush(&mut batch).await; }
        }
    }
}
}

3. Semaphore for concurrency cap (not strictly a channel, but the same role):

#![allow(unused)]
fn main() {
let sem = Arc::new(Semaphore::new(10));
for url in urls {
    let permit = sem.clone().acquire_owned().await?;
    tokio::spawn(async move {
        fetch(url).await;
        drop(permit);
    });
}
}

Closing Behavior

ChannelSender dropsReceiver drops
mpscrecv() returns None when all senders gonesend() returns SendError
oneshotSender drops → receiver gets RecvErrorReceiver drops → sender's send() returns Err(value)
broadcastAll senders gone → recv returns ClosedSender's send succeeds (no listeners)
watchSender drops → changed() returns ErrSender's send returns Err

This is how channels signal cancellation: just drop one end.

Common Mistakes

MistakeFix
Unbounded mpsc "for safety"Memory grows unboundedly under load
broadcast subscribers created lazily after sendsNew subscribers miss old messages
Cloning oneshot::Sender (it doesn't impl Clone)Use mpsc if you need multiple senders
Holding a sender forever, expecting receiver to stop on NoneDrop the sender when done so recv() returns None
watch for streaming eventsIt's latest-value only, not a stream

[!NOTE] Channel selection is a design question: what should happen when the receiver is slow? Block (mpsc bounded), drop (try_send), buffer-and-replay (broadcast), or "always show current" (watch). Pick the one that matches your domain semantics.

Interview Follow-ups

  • "Why is mpsc only single-consumer?" — Multi-consumer would need synchronization on every receive. crossbeam_channel is mpmc with that overhead; tokio's mpsc is faster by skipping it.
  • "How is async channel different from sync?" — Async sends/receives integrate with the executor — they .await instead of blocking the thread.
  • "What's Notify?" — Tokio's signal primitive — like a Condvar for async. Useful for "wake up the waiting task when something changes."