6.5 · Channels: mpsc and mpmc Message Passing
Domain 6 — Concurrency and Parallelism Duration: ~15 minutes Library components:
std::sync::mpsc,std::sync::mpmc,Sender,Receiver,SyncSender,RecvError,SendError,TryRecvError,TrySendError
Introduction
"Do not communicate by sharing memory. Instead, share memory by communicating." Channels move the ownership of values between threads. After send(value), the producer cannot use the value. The compiler does more than prevent the data race: the program cannot express it.
This tutorial shows:
mpsc::channel(unbounded): multiple producers, one consumer, and iteration on the receiver.mpsc::sync_channel(bounded): backpressure,try_send, and the rendezvous channel with a bound of 0.- The error types:
SendError,TrySendError,TryRecvError,RecvTimeoutError. Also the shutdown pattern in which you drop the sender. - The unstable
std::sync::mpmcmodule (nightly as of 1.99): cloneable receivers and the worker pool that they make possible. - How to select between channels and shared state.
mpsc::channel: Multi-Producer, Single-Consumer
Figure: mpsc Topology
channel() returns a (Sender<T>, Receiver<T>) pair with an unbounded buffer, so send never blocks. Sender is Clone: this is the "mp" (multi-producer) in the name. Receiver is not Clone: this is the "sc" (single-consumer).
// LogRecord is a struct with the fields producer: u32, seq: u32, and message: String.
let (tx, rx) = mpsc::channel::<LogRecord>();
for producer in 1..=3 {
let tx = tx.clone(); // each producer thread gets its own Sender
thread::spawn(move || {
for seq in 0..4 {
let message = format!("worker {producer} event {seq}");
// send moves the record into the channel.
tx.send(LogRecord { producer, seq, message }).expect("receiver alive");
}
}); // the clone drops when the producer thread finishes
}
drop(tx); // drop the ORIGINAL Sender too (the list below gives the reason)
let mut records: Vec<LogRecord> = Vec::new();
for record in rx { // blocks for each message, ends when the channel closes
records.push(record);
}
// records now holds all 12 records (3 producers, 4 records each)
Three details are important:
sendmoves the value. The ownership goes through the channel. There is no clone, no lock, and no aliasing.- The receiver loop ends when every
Senderis gone. If the coordinator does not drop the originaltx, the result is a common deadlock. The loop waits forever for a sender that will never send and will never disconnect. - The order is FIFO for each sender. Messages from one sender arrive in the order of the sends. The interleaving between senders depends on the scheduler. Assert the order of one sender. Never assert the interleaving between senders.
06_15_mpsc_basics.rs prints:
collected 12 records from 3 producers
per-producer FIFO order verified
total message bytes: 192
recv got: single shot
All assertions passed.
sync_channel: Bounds and Backpressure
sync_channel(bound) sets a limit on the buffer of messages in transit. When the buffer is full, send blocks until the consumer receives a message. This block is backpressure. Because of backpressure, the memory use of a bounded channel does not grow when a fast producer supplies a slow consumer:
let (tx, rx) = mpsc::sync_channel::<u32>(2); // the buffer holds 2 messages
tx.send(1).unwrap(); // buffered
tx.send(2).unwrap(); // buffered, and the buffer is now full
match tx.try_send(3) { // does not block: it returns an error immediately
Err(TrySendError::Full(rejected)) => assert_eq!(rejected, 3), // Full returns the value
_ => unreachable!(), // rx is alive and the buffer is full
}
try_send returns the rejected value in the error. Nothing disappears silently. Its second failure mode is TrySendError::Disconnected(value): the receiver is gone.
The rendezvous channel
sync_channel(0) has no buffer. Each send blocks until a recv is in progress at the same time. Thus the message transfer is a synchronization point. It is a handshake between two threads that needs no other mechanism:
let (tx_ready, rx_ready) = mpsc::sync_channel::<&str>(0); // bound 0: no buffer
// The worker thread owns tx_ready, and the main thread owns rx_ready.
// worker: tx_ready.send("initialized") blocks until main calls recv
// main: rx_ready.recv() returns Ok("initialized")
How to select a bound:
0: a handshake in which the two threads move in step.- A small bound (approximately the CPU count): the buffer absorbs bursts, and memory use has a limit.
- Unbounded
channel(): only when a different part of the program already limits the rate of the producer.
06_16_sync_channel.rs prints:
try_send on full buffer: rejected value 3 handed back
backpressured transfer: 20 items, in order
rendezvous handshake: initialized
try_send after receiver dropped: Disconnected(99)
All assertions passed.
Channel Errors: Each Way a Channel Operation Fails
Each error type answers one precise question. Empty and Disconnected need opposite reactions:
| API | Error | Meaning | Reaction |
|---|---|---|---|
send | SendError(value) | The receiver is gone | Log or requeue the returned value |
try_send | TrySendError::Full(value) | The buffer is full | Apply a backoff policy or a drop policy |
try_send | TrySendError::Disconnected(value) | The receiver is gone | Stop the producer |
try_recv | TryRecvError::Empty | No message at this time, and senders are alive | Poll again later |
try_recv | TryRecvError::Disconnected | All senders are gone, and the queue is empty | Exit the loop |
recv | RecvError | All senders are gone | Exit the loop |
recv_timeout | RecvTimeoutError::Timeout | No message during the timeout | Do housekeeping, then retry |
recv_timeout | RecvTimeoutError::Disconnected | All senders are gone (it returns immediately) | Exit the loop |
The shutdown pattern is a result of this design. Do not send a special "STOP" sentinel value. Drop the senders. The for job in rx loop of the consumer then ends. It ends even if a producer panicked, because the Sender of that producer drops during unwinding:
// tx and rx are the two ends of mpsc::channel::<u32>().
// The worker adds the jobs until the channel closes, then returns the sum.
let worker = thread::spawn(move || rx.into_iter().sum::<u32>());
for job in 1..=5 { tx.send(job).unwrap(); } // sends 1, 2, 3, 4, 5
drop(tx); // no sender remains: this IS the shutdown signal
assert_eq!(worker.join().unwrap(), 15); // 1 + 2 + 3 + 4 + 5
06_17_channel_errors.rs prints:
try_recv: Empty (senders alive, nothing queued)
try_recv: Disconnected (all senders dropped)
send failed, value recovered: "audit-event-4711"
recv_timeout with queued message: Ok("prompt delivery")
recv_timeout on silent channel: Timeout after ~10ms
recv_timeout on closed channel: Disconnected (immediate)
worker processed sum 15, then shut down cleanly
All assertions passed.
std::sync::mpmc: Cloneable Receivers (Nightly)
The mpmc module (multi-producer, multi-consumer) is still unstable as of Rust 1.99. Its feature is mpmc_channel, and its tracking issue is #126840. The main difference from mpsc is that Receiver is Clone. Thus a pool of workers can receive from one shared queue. The channel delivers each message to exactly one consumer: it is a job queue, not a broadcast.
Figure: mpmc Worker Pool
#![feature(mpmc_channel)] // nightly only
use std::sync::mpmc;
let (tx, rx) = mpmc::channel::<u64>();
for _ in 0..3 {
let rx = rx.clone(); // impossible with mpsc::Receiver!
// Each worker takes jobs from the same queue. Two workers never get the same job.
// `process` is a placeholder for the work on one job.
thread::spawn(move || for job in rx { process(job); });
}
By design, the API matches the mpsc API: channel, sync_channel, and the same method names. A migration will be mostly a change of the use line. Until mpmc is stable, stable Rust has these alternatives:
Arc<Mutex<mpsc::Receiver<T>>>.- One channel for each worker, with round-robin dispatch.
- The crossbeam-channel crate, which was the model for this API.
The example binary is nightly-gated. Run it with this command:
cargo +nightly run -p domain-06-concurrency --features nightly --bin 06_18_mpmc_channel
06_18_mpmc_channel.rs prints this output. The number of jobs for each worker changes between runs:
worker drained 2 jobs (varies)
worker drained 15 jobs (varies)
worker drained 13 jobs (varies)
total jobs: 30, total sum: 465
All assertions passed.
Choosing: Channels or Shared State
| Shape of the problem | Tool |
|---|---|
| Data flows in one direction, and the stages make a pipeline | Channels |
| Work items go to a pool of workers | mpmc (nightly) or Arc<Mutex<Receiver>> |
| Many threads read and update one long-lived structure | Arc<Mutex<T>> or Arc<RwLock<T>> (6.2) |
| One value with single-step updates | Atomics (6.3) |
| All threads must reach a checkpoint together | Barrier (6.7) |
Channels are the correct tool when the transfer of ownership agrees with the problem domain: jobs, log records, results. They remove aliasing by construction. The cost is the move of the data and a queue overhead for each message.
Shared state is the correct tool when the threads really use one common structure, for example a cache or a metrics registry. For such a structure, a copy through a channel would be artificial.
You can use the two together. A common architecture uses channels to distribute the work and one Arc<Mutex<_>> for the final aggregate.
Summary
| Concept | Key point |
|---|---|
mpsc::channel | Unbounded. send never blocks. Sender is Clone, and Receiver is not. |
send(value) | Moves the ownership. The program cannot express a data race on the value. |
| Receiver iteration | for msg in rx ends when every sender is gone |
Drop the original tx | If you do not, the consumer loop never ends |
| FIFO for each sender | Guaranteed. The interleaving between senders is not guaranteed. |
sync_channel(n) | Bounded. A full buffer blocks send: this is backpressure. |
sync_channel(0) | Rendezvous: each send waits for a matching recv |
try_send / try_recv | They do not block. A try_send error returns the value. A try_recv error is Empty or Disconnected. |
recv_timeout | Blocks for a limited time. Disconnected returns immediately. |
| Shutdown pattern | Drop the senders. Do not send sentinel messages. |
std::sync::mpmc | Nightly (1.99): cloneable receivers, and the channel delivers each message exactly once |
Code Examples
| File | Description |
|---|---|
06_15_mpsc_basics.rs | Cloned senders, ownership transfer, receiver iteration, FIFO checks |
06_16_sync_channel.rs | Bounds, try_send/Full, backpressure, rendezvous channel |
06_17_channel_errors.rs | All the error types and the drop-the-sender shutdown pattern |
06_18_mpmc_channel.rs | (nightly) Cloneable receivers: 2 producers and a pool of 3 consumers |