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

为何tokio::select!高频选中某分支?异步通道超时异常排查

问题分析:Tokio select! 分支执行不符合预期

代码复现

use tokio::sync::{mpsc, oneshot};
use std::time::Duration;

#[derive(Debug)]
pub struct Message{
    content:String,
    id:i32
}

impl Message {
    pub fn new(s : String,id:i32) -> Message{
        Message { 
            content:s,
            id 
        }
    }
}


async fn msg_stream(sender : mpsc::Sender<Message>) {
    loop{
        tokio::time::sleep(Duration::from_secs(1)).await;
        let m = Message::new("abc".to_string(),1);
        println!("message = {:?}",m);
        if let Err(e) = sender.send(m).await{
            println!("channel is closed,{}",e);
            break
        }
    }
}

async fn read_stream(mut receiver : mpsc::Receiver<Message>){
    let (tx, mut rx) = oneshot::channel::<()>();
    loop{
        tokio::select! {
            Err(_) = tokio::time::timeout(Duration::from_secs(3),& mut rx) => {
                println!("time has elapsed");
                break
            }
            message = receiver.recv() =>{
                println!("was receiver message = {:?} ",message)
            }
        }
    }
    println!("END OF STREAM");

}


#[tokio::main]
async fn main() {
    let (tx,rx) = mpsc::channel::<Message>(8);
    tokio::join!(msg_stream(tx),read_stream(rx));
    println!("end of the programm");
}

现象对比

  • 消息每秒发送一次:程序无限循环接收消息,超时分支从未触发,持续打印消息发送与接收日志。
  • 消息每3秒发送一次:3秒超时触发,读取流退出,随后发送方因通道关闭终止,程序结束。

问题原因

核心问题出在read_stream的超时分支设计:

  1. 你创建的oneshot通道从未调用tx.send(()),导致rx永远处于等待状态。
  2. 每次循环中的tokio::time::timeout(Duration::from_secs(3), &mut rx)都会重新启动3秒倒计时,而非维持连续的超时计时。
  3. 当消息每秒发送时,receiver.recv()分支1秒后就就绪,select!会优先处理就绪分支,超时倒计时还未走完就进入下一次循环重新计时,因此永远无法触发超时。
  4. 当消息每3秒发送时,第一次循环中3秒超时先于消息接收就绪,所以超时分支被触发,读取流退出。

tokio::select!的分支选择逻辑并非随机,优先执行已就绪的分支;仅当多个分支同时就绪时才会随机选择。这里不存在超时消息丢失,只是超时逻辑设计错误。

解决方案

要实现“3秒未收到消息则退出”的逻辑,直接用tokio::time::sleep作为超时分支即可,无需依赖未使用的oneshot通道:

修改后的read_stream函数:

async fn read_stream(mut receiver : mpsc::Receiver<Message>){
    loop{
        tokio::select! {
            _ = tokio::time::sleep(Duration::from_secs(3)) => {
                println!("time has elapsed");
                break
            }
            message = receiver.recv() =>{
                match message {
                    Some(msg) => println!("was receiver message = {:?} ", msg),
                    None => {
                        println!("sender closed channel");
                        break;
                    }
                }
            }
        }
    }
    println!("END OF STREAM");
}

修改后每次循环会同时等待“3秒睡眠”和“消息接收”:

  • 3秒内无消息则触发超时退出;
  • 3秒内收到消息则处理后进入下一次循环,重新开始3秒计时。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 11:46:01