8.2 消息传递:channel传票


8.2 消息传递:channel 传票

本节摘要: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 { ... } 的迭代形态。

多生产者:传票窗口clone

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 多生产者单接收者,发送端 clone 即增窗口;
  • 接收端可迭代,关闭自动收尾;
  • 流水线与任务分发是通道的主辖区,下一节处理"必须共享"的场合。

传票的产权语义

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,接收方先死时发送失败要处理。这两个设计都源自同一法理——通道状态由两端的产权共同决定,没有隐藏的全局开关。

try_recv 与轮询的正确姿势

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"的对照表,九成问题当场显形。

select 与多路等待的缺口

标准库 mpsc 没有内置"同时等多个通道"的原语,这是初学者常撞的墙。三条出路:一是架构上合并通道——多个上游 clone 同一个 sender,接收侧单通道分派,多数业务其实只需要这一条;二是引入带 select 的运行时或库(异步运行时原生支持);三是用轮询加超时的 try_recv 组合,仅限已有事件循环的宿主。第一路最被低估:把"等谁"的问题转化成"消息里带谁",通道模型的美德恰恰是让等待方无需关心来源,加个枚举消息类型就够。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U