Skip to content

Ch 35: std::sync::mpsc - 消息通道

mpsc(Multi-Producer, Single-Consumer)提供了线程间的消息传递机制。一个通道可以有多个发送者,但只能有一个接收者。

use std::sync::mpsc;
use std::thread;
fn main() {
// 创建通道
let (tx, rx) = mpsc::channel();
// 发送消息
thread::spawn(move || {
tx.send(42).unwrap();
});
// 接收消息(阻塞直到收到)
let result = rx.recv().unwrap();
println!("收到: {}", result);
}
// Sender
impl<T> Sender<T> {
pub fn send(&self, t: T) -> Result<(), SendError<T>>
pub fn try_send(&self, t: T) -> Result<(), TrySendError<T>>
pub fn is_disconnected(&self) -> bool
pub fn capacity(&self) -> Option<usize>
}
// Receiver
impl<T> Receiver<T> {
pub fn recv(&self) -> Result<T, RecvError>
pub fn try_recv(&self) -> Result<T, TryRecvError>
pub fn iter(&self) -> Iter
pub fn try_iter(&self) -> TryIter
}

多个线程可以同时发送:

use std::sync::mpsc;
use std::thread;
fn main() {
let (tx, rx) = mpsc::channel();
// 创建多个生产者
for i in 0..3 {
let tx = tx.clone();
thread::spawn(move || {
tx.send(format!("来自线程 {}", i)).unwrap();
});
}
// 接收所有消息
drop(tx); // 显式drop主发送者
for msg in rx {
println!("收到: {}", msg);
}
}

sync_channel是带缓冲区的同步通道:

use std::sync::mpsc;
use std::thread;
fn main() {
// 创建容量为2的同步通道
let (tx, rx) = mpsc::sync_channel(2);
// 发送者会阻塞直到缓冲区有空间
tx.send(1).unwrap();
tx.send(2).unwrap();
// 容量已满,send会阻塞
// tx.send(3).unwrap(); // 会阻塞
println!("通道容量: {:?}", tx.capacity());
// 接收消息
println!("收到: {}", rx.recv().unwrap());
println!("收到: {}", rx.recv().unwrap());
}

方法签名:

pub fn sync_channel<T>(cap: usize) -> (Sender<T>, Receiver<T>)
use std::sync::mpsc;
use std::thread;
fn main() {
let (tx, rx) = mpsc::channel();
// try_send立即返回(不阻塞)
match tx.try_send(1) {
Ok(_) => println!("发送成功"),
Err(_) => println!("通道已满"),
}
// try_recv立即返回(不阻塞)
match rx.try_recv() {
Ok(msg) => println!("收到: {}", msg),
Err(_) => println!("通道为空"),
}
}
use std::sync::mpsc;
use std::thread;
fn main() {
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
for i in 0..5 {
tx.send(i).unwrap();
}
});
// for循环迭代器(阻塞直到通道关闭)
for msg in rx.iter() {
println!("收到: {}", msg);
}
}

迭代器类型:

  • iter() - 阻塞迭代
  • try_iter() - 非阻塞迭代

8. select! 宏(需要启用crossbeam或mpsc扩展)

Section titled “8. select! 宏(需要启用crossbeam或mpsc扩展)”

标准库mpsc不支持select,但crossbeam crate提供:

// Cargo.toml添加
// crossbeam-channel = "0.5"
use crossbeam_channel::{select, bounded};
fn main() {
let (tx1, rx1) = bounded(1);
let (tx2, rx2) = bounded(1);
tx1.send("A").unwrap();
tx2.send("B").unwrap();
select! {
recv(rx1, msg) => println!("收到1: {:?}", msg),
recv(rx2, msg) => println!("收到2: {:?}", msg),
default => println!("都没有"),
}
}
use std::sync::mpsc;
use std::time::{Duration, Instant};
fn main() {
let (tx, rx) = mpsc::channel();
// 使用thread::sleep模拟超时
let start = Instant::now();
let timeout = Duration::from_millis(100);
thread::spawn(move || {
thread::sleep(Duration::from_secs(1));
tx.send(42).unwrap();
});
// 手动实现超时
let result = loop {
if let Ok(msg) = rx.try_recv() {
break Ok(msg);
}
if start.elapsed() > timeout {
break Err("超时");
}
thread::sleep(Duration::from_millis(10));
};
match result {
Ok(msg) => println!("收到: {}", msg),
Err(_) => println!("接收超时"),
}
}
use std::sync::mpsc;
use std::thread;
use std::time::Duration;
struct Work {
id: u32,
data: String,
}
fn main() {
let (tx, rx) = mpsc::channel();
// 创建生产者
let producer_tx = tx.clone();
let producer = thread::spawn(move || {
for i in 0..10 {
let work = Work {
id: i,
data: format!("任务{}", i),
};
producer_tx.send(work).unwrap();
thread::sleep(Duration::from_millis(50));
}
});
// 创建多个消费者
let mut consumers = vec![];
for consumer_id in 0..2 {
let rx = rx.clone();
let handle = thread::spawn(move || {
while let Ok(work) = rx.recv() {
println!("消费者{} 处理: {} - {}", consumer_id, work.id, work.data);
thread::sleep(Duration::from_millis(100));
}
});
consumers.push(handle);
}
// 等待生产者完成
producer.join().unwrap();
// 等待消费者处理完所有消息(使用wait)
drop(tx);
for handle in consumers {
handle.join().unwrap();
}
println!("所有任务完成");
}
  1. Sender克隆:channel()的Sender可以克隆(多生产者)
  2. Receiver不可克隆:Receiver是单消费者的,不能克隆
  3. 通道关闭:所有Sender被drop后,Receiver会收到错误
  4. 同步vs异步:channel()异步(缓冲),sync_channel()同步(阻塞)
  5. 容量:sync_channel(cap)容量为0时,完全同步

std::sync::mpsc核心API:

  • channel() - 创建异步通道(无界缓冲)
  • sync_channel(cap) - 创建同步通道
  • Sender::send() - 发送(阻塞)
  • Sender::try_send() - 尝试发送(非阻塞)
  • Receiver::recv() - 接收(阻塞)
  • Receiver::try_recv() - 尝试接收(非阻塞)
  • rx.iter() - 迭代器接口

消息通道是Rust并发编程的核心模式之一,适合:

  • 生产者-消费者场景
  • 任务分发
  • 线程间通信