阶段四 · 并发编排与 IO

异步消息传递

一句话总结

Tokio 提供四种异步 channel:mpsc(多生产者单消费者,最常用),oneshot(一次性传一个值,常用于返回结果),broadcast(广播给所有订阅者),watch(只关心最新值).
它们对应 Go channel 的不同使用模式,但分工更明确.

为什么优先用消息传递

上一篇用锁共享状态,本篇换一种并发协作思路:任务之间不共享内存,而是互相发消息.这正是 Go 那句名言的主张--"不要通过共享内存来通信,而要通过通信来共享内存".

消息传递的好处是把"谁能改数据"收敛到单个任务,天然避免锁竞争.Tokio 没有像 Go 那样只给一种 channel,而是按使用模式拆成四种,各司其职.

mpsc:多生产者,单消费者

最常用的一种.多个任务往里发,一个任务收--典型的"工作分发"或"事件汇聚"模型:

use tokio::sync::mpsc;

#[tokio::main]
async fn main() {
let (tx, mut rx) = mpsc::channel::<i32>(32); // 32 是缓冲容量

for i in 0..3 {
let tx = tx.clone();
tokio::spawn(async move {
tx.send(i).await.unwrap(); // 缓冲满时 .await 等待
});
}
drop(tx); // 丢弃最初的 tx,否则 rx 永远等不到关闭

while let Some(v) = rx.recv().await {
println!("收到 {v}");
}
}

mpsc ≈ Go 的带缓冲 channel + 多个 goroutine 写,一个 goroutine 读.
区别在于:Go 的 channel 读写两端不分家;Tokio 把发送端(Sender,可克隆)和接收端(Receiver,唯一)类型上分开,"多生产者单消费者"由类型保证,而非靠约定.

两个新手必踩的点

  • 记得 drop 发送端:只要还有任何 Sender 存活,rx.recv() 就认为"可能还有消息",不会返回 None.上例中若不 drop(tx),while let 会永远卡住.
  • recv 返回 Option:所有发送端都关闭后,recv() 返回 None,这正是循环的退出信号--对应 Go 中 v, ok := <-chok == false.

oneshot:一次性返回一个值

顾名思义,只能发送一次,接收一次.最常见的用途是"给一个任务派活,并等它把单个结果送回来":

use tokio::sync::oneshot;

let (tx, rx) = oneshot::channel();

tokio::spawn(async move {
let result = compute().await;
let _ = tx.send(result); // 只能 send 一次,消耗 tx
});

let value = rx.await.unwrap(); // 接收端本身就是 Future,直接 await

oneshot一张回执单:派单时附上它,对方办完把结果填上送回,单子用一次就作废.
它常和 mpsc 搭配--请求里带一个 oneshot 的 tx,处理方用它回传结果,实现"请求-响应"模式.

broadcast:广播给所有接收者

多生产者,多消费者,且每条消息每个接收者都能收到一份(区别于 mpsc 的"一条消息只有一个消费者拿到"):

use tokio::sync::broadcast;

let (tx, mut rx1) = broadcast::channel(16);
let mut rx2 = tx.subscribe(); // 再订阅一个接收端

tx.send("公告").unwrap();
assert_eq!(rx1.recv().await.unwrap(), "公告");
assert_eq!(rx2.recv().await.unwrap(), "公告"); // 两个接收端各收到一份

适用于聊天室广播,配置变更通知等"一对多"场景.接收太慢,缓冲被覆盖时,会收到 Lagged 错误,提示丢了若干条旧消息.

watch:只关心最新值

特化的广播:接收者不在乎中间过程,只要最新状态.channel 里只保留最新一个值,新值覆盖旧值:

use tokio::sync::watch;

let (tx, mut rx) = watch::channel("初始配置");

tokio::spawn(async move {
while rx.changed().await.is_ok() {
println!("配置更新为:{}", *rx.borrow()); // 读当前最新值
}
});

tx.send("新配置").unwrap();

典型用途:配置热更新,共享一个"当前状态",传播关闭信号(第 13 篇会用它做优雅关闭).

四种 channel 选型对照

channel 生产者 消费者 消息分发 典型场景
mpsc 每条被一个消费者取走 任务队列,事件汇聚
oneshot 单(一次) 单(一次) 仅一个值 返回单个结果,请求-响应
broadcast 每条每个消费者各一份 广播通知,聊天室
watch 单/多 只保留最新值 配置热更新,状态共享

与 Go channel 的关系

Go 用一种 channel(配合方向标注,缓冲,close)覆盖所有这些模式,灵活但语义靠约定.Tokio 把高频模式做成四种专门类型:代码意图一眼可辨(看到 oneshot 就知道是单次返回),也减少了误用.

多数业务里 mpsc + oneshot 的组合就够用了.

快速回顾

  • 思路:用消息传递替代共享锁,把"谁能改数据"收敛到单个任务.
  • mpsc:多生产者单消费者,最常用;注意 drop 所有 Sender 才会让 recv 返回 None.
  • oneshot:一次性单值,接收端即 Future;常与 mpsc 配合实现请求-响应.
  • broadcast:一对多,每条消息每个订阅者各收一份,接收慢会 Lagged 丢消息.
  • watch:只保留最新值,适合配置热更新,状态/关闭信号传播.
  • 对照 Go:Go 一种 channel 打天下,Tokio 拆成四种专用类型,意图更清晰.

动手练习

  1. mpsc 基础:用 mpsc 实现 3 个生产者各发 5 个数字,1 个消费者全部收齐后打印总数,注意正确关闭.
  2. 请求-响应:用 mpsc + oneshot 实现请求-响应:消费者收到请求后,通过请求里携带的 oneshot 回传结果.
  3. broadcast 实验:用 broadcast 让两个接收者都收到同一条消息,并制造一次 Lagged 观察其表现.
  4. watch 热更新:用 watch 模拟配置热更新:一个任务持续读最新配置,主任务多次更新.
  5. 选型说明:给上面四种场景各写一句话,说明"为什么这里该用这种 channel 而非别的".