为何Tokio Notify示例中的Channel是MPSC?求MPMC失效测试代码
问题:为什么Tokio Notify示例中的Channel是MPSC而非MPMC?
我在学习Tokio的同步原语时,看到Notify的示例代码里定义了一个Channel<T>,注释明确说明它是**单消费者多生产者(MPSC)**的,但我一开始没搞懂为什么它不能作为多消费者多生产者(MPMC)使用。
原示例代码如下:
use tokio::sync::Notify; use std::collections::VecDeque; use std::sync::Mutex; struct Channel<T> { values: Mutex<VecDeque<T>>, notify: Notify, } impl<T> Channel<T> { pub fn send(&self, value: T) { self.values.lock().unwrap() .push_back(value); // Notify the consumer a value is available self.notify.notify_one(); } // This is a single-consumer channel, so several concurrent calls to // `recv` are not allowed. pub async fn recv(&self) -> T { loop { // Drain values if let Some(value) = self.values.lock().unwrap().pop_front() { return value; } // Wait for values to be available self.notify.notified().await; } } }
这个Channel的核心逻辑是:
- 若
values队列中有元素,消费者调用recv会直接取走队首元素并返回 - 若队列为空,消费者会挂起,直到生产者调用
send时发出通知
我自己尝试写测试代码时没发现问题,直到用下面的测试代码才复现了它作为MPMC使用时的不安全之处:
use std::sync::Arc; #[tokio::main] async fn main() { let mut i = 0; loop{ let ch = Arc::new(Channel { values: Mutex::new(VecDeque::new()), notify: Notify::new(), }); let mut handles = vec![]; for i in 0..100{ if i % 2 == 1{ for _ in 0..2{ let sender = ch.clone(); tokio::spawn(async move{ sender.send(1); }); } }else{ for _ in 0..2{ let receiver = ch.clone(); let handle = tokio::spawn(async move{ receiver.recv().await; }); handles.push(handle); } } } futures::future::join_all(handles).await; i += 1; println!("No.{i} loop finished."); } }
问题复现说明
如果程序卡在某个循环中无法进入下一轮,就说明存在未完成的消费者任务——本质是通知丢失了。
导致这个问题的核心原因是:
当多个消费者都挂在notify.notified().await上时,生产者调用notify_one()只会唤醒其中一个消费者。如果被唤醒的消费者拿到锁后发现队列已经被其他消费者取空,它会再次挂起;但如果此时队列中还有新元素(比如唤醒期间又有生产者发送了元素,但对应的notify_one()已经调用过),后续的消费者就可能永远等不到通知,导致任务卡住。
举个具体的场景:
- 两个消费者都发现队列为空,进入挂起状态
- 生产者放入一个元素,调用
notify_one()唤醒消费者A - 消费者A还未拿到锁时,另一个生产者又放入一个元素,调用
notify_one()唤醒消费者B - 消费者B先拿到锁,取走了两个元素后返回
- 消费者A拿到锁,发现队列空了,再次挂起
- 此时没有更多生产者发送元素,消费者A会永远挂起,导致
join_all无法完成,程序卡住
内容的提问来源于stack exchange,提问作者spongecaptain
相关产品推荐
相关产品推荐

