You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

为何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()已经调用过),后续的消费者就可能永远等不到通知,导致任务卡住。

举个具体的场景:

  1. 两个消费者都发现队列为空,进入挂起状态
  2. 生产者放入一个元素,调用notify_one()唤醒消费者A
  3. 消费者A还未拿到锁时,另一个生产者又放入一个元素,调用notify_one()唤醒消费者B
  4. 消费者B先拿到锁,取走了两个元素后返回
  5. 消费者A拿到锁,发现队列空了,再次挂起
  6. 此时没有更多生产者发送元素,消费者A会永远挂起,导致join_all无法完成,程序卡住

内容的提问来源于stack exchange,提问作者spongecaptain

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.07 07:15:34