本节摘要:channel 把"共享内存"换成"传递所有权"——消息发出去,发送方就失去产权,接收方独占处理,竞争在制度上不存在。本节讲 mpsc 通道的建立、多生产者形态与所有权语义。读完你能用通道搭建流水线式并发。
use std::sync::mpsc; use std::thread; let (tx, rx) = mpsc::channel(); thread::spawn(move || { let verdicts = vec!["无罪", "缓刑", "罚金"]; for v in verdicts { tx.send(v).unwrap(); // 发出即失去产权 } }); for _ in 0..3 { println!("送达:{}", rx.recv().unwrap()); // 阻塞等待传票 }
send 之后那份数据在本线程已不可用——不是纪律要求,是类型判决。谁持有、谁处理、谁释放,全程唯一。若发送端全部关闭,recv 返回 Err,循环自然终止,可写成 for msg in rx { ... } 的迭代形态。
let (tx, rx) = mpsc::channel(); let tx2 = tx.clone(); // 发送端可复制,接收端唯一 thread::spawn(move || { tx.send("窗口一").unwrap(); }); thread::spawn(move || { tx2.send("窗口二").unwrap(); }); println!("{}", rx.recv().unwrap()); // 两张传票先到先得

mpsc 通道传递的是产权本身,这句话的所有推论都值得写下来:发送后原值作废、接收方独占、通道关闭由 drop 自动声明。
use std::sync::mpsc; use std::thread; fn main() { let (tx, rx) = mpsc::channel(); for i in 1..=3 { let tx = tx.clone(); // 多发送方:克隆的是发件权 thread::spawn(move || { tx.send(format!("卷宗{}已勘验", i)).expect("接收方存活"); }); } drop(tx); // 主线程交出发件权,否则循环不结束 for msg in rx { // 迭代到全部发送方 drop 为止 println!("收到 {}", msg); } }
两处细节是初学高频坑:不 drop 原始 tx,for msg in rx 永远等不到关闭;send 返回 Result,接收方先死时发送失败要处理。这两个设计都源自同一法理——通道状态由两端的产权共同决定,没有隐藏的全局开关。
use std::sync::mpsc; use std::time::Duration; fn main() { let (tx, rx) = mpsc::channel(); thread::spawn(move || { thread::sleep(Duration::from_millis(50)); tx.send("迟到的证词").unwrap(); }); loop { match rx.try_recv() { Ok(msg) => { println!("{}", msg); break; } Err(mpsc::TryRecvError::Empty) => { thread::sleep(Duration::from_millis(10)); // 让出 CPU } Err(mpsc::TryRecvError::Disconnected) => break, } } }
忙等不 sleep 会烧满一个核,try_recv 轮询只该出现在已有事件循环的宿主里;纯消费者直接用阻塞 recv 或迭代器。
sync_channel(N) 设定缓冲上限,满了之后 send 阻塞——上游自动减速,这就是背压,不需要任何额外机制。
use std::sync::mpsc; use std::thread; fn main() { let (tx, rx) = mpsc::sync_channel(2); // 缓冲两个传票 let producer = thread::spawn(move || { for i in 0..5 { tx.send(i).unwrap(); // 第 3 个起阻塞,等消费者取件 println!("送出 {}", i); } }); let consumer = thread::spawn(move || { for _ in 0..5 { thread::sleep(std::time::Duration::from_millis(20)); println!("取件 {:?}", rx.recv()); } }); producer.join().unwrap(); consumer.join().unwrap(); }
无界 channel 在生产快于消费时内存无上限增长,服务端代码默认应选有界版,容量数值本身也是一条可调参数。
| 形态 | 缓冲 | send 行为 | 场景 |
|---|---|---|---|
| channel() | 无界 | 立即返回 | 吞吐优先、消费可靠 |
| sync_channel(0) | 零 | 阻塞至接收方就位 | 严格握手、逐件交接 |
| sync_channel(N) | 有界 | 满则阻塞 | 背压、保护内存 |
零缓冲版本常被忽视却极有价值:它把通道变成同步会面点,发送方确知"对方已接到手"才继续,等价于一次零共享的状态交接。排查通道相关卡死时,先画出"谁阻塞在 send、谁阻塞在 recv"的对照表,九成问题当场显形。
标准库 mpsc 没有内置"同时等多个通道"的原语,这是初学者常撞的墙。三条出路:一是架构上合并通道——多个上游 clone 同一个 sender,接收侧单通道分派,多数业务其实只需要这一条;二是引入带 select 的运行时或库(异步运行时原生支持);三是用轮询加超时的 try_recv 组合,仅限已有事件循环的宿主。第一路最被低估:把"等谁"的问题转化成"消息里带谁",通道模型的美德恰恰是让等待方无需关心来源,加个枚举消息类型就够。