Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

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::mpmc module (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:

  • send moves the value. The ownership goes through the channel. There is no clone, no lock, and no aliasing.
  • The receiver loop ends when every Sender is gone. If the coordinator does not drop the original tx, 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:

APIErrorMeaningReaction
sendSendError(value)The receiver is goneLog or requeue the returned value
try_sendTrySendError::Full(value)The buffer is fullApply a backoff policy or a drop policy
try_sendTrySendError::Disconnected(value)The receiver is goneStop the producer
try_recvTryRecvError::EmptyNo message at this time, and senders are alivePoll again later
try_recvTryRecvError::DisconnectedAll senders are gone, and the queue is emptyExit the loop
recvRecvErrorAll senders are goneExit the loop
recv_timeoutRecvTimeoutError::TimeoutNo message during the timeoutDo housekeeping, then retry
recv_timeoutRecvTimeoutError::DisconnectedAll 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 problemTool
Data flows in one direction, and the stages make a pipelineChannels
Work items go to a pool of workersmpmc (nightly) or Arc<Mutex<Receiver>>
Many threads read and update one long-lived structureArc<Mutex<T>> or Arc<RwLock<T>> (6.2)
One value with single-step updatesAtomics (6.3)
All threads must reach a checkpoint togetherBarrier (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

ConceptKey point
mpsc::channelUnbounded. 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 iterationfor msg in rx ends when every sender is gone
Drop the original txIf you do not, the consumer loop never ends
FIFO for each senderGuaranteed. 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_recvThey do not block. A try_send error returns the value. A try_recv error is Empty or Disconnected.
recv_timeoutBlocks for a limited time. Disconnected returns immediately.
Shutdown patternDrop the senders. Do not send sentinel messages.
std::sync::mpmcNightly (1.99): cloneable receivers, and the channel delivers each message exactly once

Code Examples

FileDescription
06_15_mpsc_basics.rsCloned senders, ownership transfer, receiver iteration, FIFO checks
06_16_sync_channel.rsBounds, try_send/Full, backpressure, rendezvous channel
06_17_channel_errors.rsAll 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