mpsc 是 multiple producer, single consumer 的缩写,表示多个发送端、一个接收端。它是一种常见的消息通信模型:多个地方可以把消息发送到同一个 channel 中,再由一个接收方按顺序取出并处理。
前言
写服务端时,很多模块都可能需要给同一个客户端发消息。比如登录回包、消息广播、系统通知,甚至定时任务产生的推送。如果这些模块都直接操作连接对象,代码很容易碰到借用和并发访问的问题,连接本身也会变得难以维护。
一种更自然的做法是,让连接只由一个地方持有,其他模块只负责投递消息,再由持有连接的任务统一写出。mpsc 适合解决的,正是这类“多处发送、单点处理”的通信问题。
标准库的 mpsc
先看一个简单的例子:
1 | use std::sync::mpsc; |
上述例子 tx 是发送端,负责将消息放入 channel;rx 是接收端,负责从 channel 中取出消息。
现在只需要对其稍加改造,就可以变成多发送端:
1 | // ... |
tx.clone() 得到的是一个新的发送端,它与原始的 tx 指向同一个 channel。它们分别被移动到不同的线程中,两个子线程的执行顺序由系统调度决定,因此无法保证哪一条消息先到达接收端。
这里需要注意 drop(tx)。如果不释放它,原始的 tx 仍然存活于主线程,rx 会认为后续仍可能有新的消息到来,因此在队列为空后继续等待,for 循环也不会结束。
Tokio 中的 mpsc
前面的示例使用的是标准库提供的 std::sync::mpsc。它的 recv() 会阻塞当前线程,适合普通的多线程程序。
当前重写的服务端运行在 Tokio 异步运行时中,因此更多使用 tokio::sync::mpsc。当没有消息时,recv().await 会挂起当前异步任务,而不会阻塞运行时线程,其他任务仍可以继续执行。
1 | use tokio::sync::mpsc; |
与标准库的 mpsc::channel() 一样,Tokio 的 mpsc::channel() 也会创建一对发送端和接收端。不同之处在于,发送和接收都可能需要等待,因此使用 await:
1 | tx.send(message).await; |
recv() 需要可变借用接收端,因此 rx 必须声明为 mut。当 channel 关闭且其中的消息已全部被接收后,rx.recv().await 会返回 None。
最大的不同在于消息队列:mpsc::channel(32) 中的 32 是队列容量。当队列满时,tx.send(message).await 会等待接收端取走消息后再继续发送。这种设计通过背压限制队列长度,避免未处理消息在 channel 中无限堆积。
在服务端中的应用实例
前面通过 mpsc 将消息的发送和接收分开处理。在当前项目中,客户端数据由外层任务负责读取和处理;需要向客户端发送消息时,只需通过发送端投递即可。
定义 Agent 用于统一处理待发送消息:
1 | use tokio::sync::mpsc; |
与前面的 mpsc::channel(32) 相比,unbounded_channel() 不限制消息队列的容量,发送消息时也不需要 .await。
最后
实际业务中的 Agent 往往还需要处理连接状态、用户信息和协议编解码等问题;这里仅关注它通过 mpsc 统一处理出站消息的部分。
至于 Tokio 异步运行时如何驱动 Future、任务如何被唤醒并重新调度,后续再继续学习(埋坑233)。