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

如何避免嵌套tokio::select!宏中Future取消致消息丢失?

问题描述

我正在用Rust的tokio::select!宏实现一个可与多方收发消息的ConnectionsManager结构体,核心代码如下:

// 假设Message类型已定义
struct Message(String);

pub enum Destination {
    Alice,
    Bob
}

struct ConnectionsManager {
    rx_from_alice: UnboundedReceiver<Message>,
    tx_to_alice:  UnboundedSender<Message>,
    rx_from_bob: UnboundedReceiver<Message>,
    tx_to_bob:  UnboundedSender<Message>,
}

impl ConnectionsManager  {
    async fn listen(&mut self) -> Result<Message, anyhow::Error> {
        select! {
            received_message = self.rx_from_alice.next() => {
                return self.process_received_message(received_message).await
            }
            received_message = self.rx_from_bob.next() => {
                return self.process_received_message(received_message).await
            }
        }
    }

    async fn send(&mut self, message: Message, dest: Destination) -> Result<(), anyhow::Error> {
        match dest {
            Destination::Alice => self.tx_to_alice.send(message).await,
            Destination::Bob => self.tx_to_bob.send(message).await,
        }
    }

    async fn process_received_message(&mut self, received_message: Option<Message>) -> Result<Message, anyhow::Error> {
        // 消息处理逻辑
        unimplemented!()
    }
}

原本基于mpsc通道的Stream实现不会丢失消息,但当调用ConnectionsManager::listen()的代码也使用select!时,出现了消息丢失风险:

async fn listen_to_a_lot_of_stuff(&mut self) -> Result<Message, anyhow::Error> {
    let mut connections_manager = ConnectionsManager::new(/* 初始化参数 */);
    loop {
        select! {
            processed_message = connections_manager.listen() => {
                // 处理返回结果
            }
            some_other_result = some_other_future => {
                // 处理其他逻辑
            }
        }
    }
}

如果some_other_future先完成,ConnectionsManager::listen()可能已经从通道接收了消息,但还没完成process_received_message(),此时Future被取消会导致消息丢失。

此前尝试的方案均失败:

  • 保存connections_manager.listen()的Future或转为Task不可行:需要在不等待该Future时调用ConnectionsManager::send(),持有Future会导致无法获取可变引用调用send方法。
  • 将TX通道交给Alice和Bob,仅保留单个RX通道:无法区分是哪一方的TX通道关闭,单RX通道返回None时无法判断是全部关闭还是某一方关闭。
解决方案

核心方案:后台任务+内部消息队列

把消息接收、处理逻辑放到独立的Tokio后台任务中,该任务持续运行并处理所有来自Alice和Bob的消息,处理完成后将结果发送到内部通道。外层代码仅监听内部通道的结果,这样即使外层select!切换到其他Future,后台任务仍会继续处理消息,完全避免丢失。同时可以在后台任务中追踪每个连接的关闭状态。

修正后的完整代码示例

use tokio::sync::{mpsc, unbounded_channel, UnboundedReceiver, UnboundedSender};
use anyhow::{Result, anyhow};

#[derive(Debug, Clone)]
struct Message(String);

#[derive(Debug, Clone, Copy)]
pub enum Destination {
    Alice,
    Bob,
}

struct ConnectionsManager {
    tx_to_alice: UnboundedSender<Message>,
    tx_to_bob: UnboundedSender<Message>,
    rx_processed: mpsc::Receiver<Result<Message>>,
    shutdown_tx: mpsc::Sender<()>,
}

impl ConnectionsManager {
    pub fn new() -> Self {
        // 创建与Alice、Bob的通信通道
        let (tx_to_alice, rx_from_alice) = unbounded_channel();
        let (tx_to_bob, rx_from_bob) = unbounded_channel();
        // 创建内部处理结果通道(有界通道防止内存溢出)
        let (tx_processed, rx_processed) = mpsc::channel(32);
        // 创建后台任务关闭信号
        let (shutdown_tx, mut shutdown_rx) = mpsc::channel(1);

        // 启动后台消息处理任务
        tokio::spawn(async move {
            let mut rx_from_alice = rx_from_alice;
            let mut rx_from_bob = rx_from_bob;

            loop {
                tokio::select! {
                    // 处理Alice的消息
                    msg = rx_from_alice.recv() => {
                        match msg {
                            Some(message) => {
                                let result = Self::process_received_message(message).await;
                                // 发送处理结果到内部通道,忽略发送失败(如外层已关闭)
                                let _ = tx_processed.send(result).await;
                            }
                            None => {
                                println!("Alice连接已关闭");
                                // 若Bob也关闭则退出任务
                                if rx_from_bob.is_closed() {
                                    break;
                                }
                                // 后续仅监听Bob的消息
                                loop {
                                    tokio::select! {
                                        msg = rx_from_bob.recv() => {
                                            match msg {
                                                Some(message) => {
                                                    let result = Self::process_received_message(message).await;
                                                    let _ = tx_processed.send(result).await;
                                                }
                                                None => break,
                                            }
                                        }
                                        _ = shutdown_rx.recv() => break,
                                    }
                                }
                                break;
                            }
                        }
                    }
                    // 处理Bob的消息
                    msg = rx_from_bob.recv() => {
                        match msg {
                            Some(message) => {
                                let result = Self::process_received_message(message).await;
                                let _ = tx_processed.send(result).await;
                            }
                            None => {
                                println!("Bob连接已关闭");
                                if rx_from_alice.is_closed() {
                                    break;
                                }
                                // 后续仅监听Alice的消息
                                loop {
                                    tokio::select! {
                                        msg = rx_from_alice.recv() => {
                                            match msg {
                                                Some(message) => {
                                                    let result = Self::process_received_message(message).await;
                                                    let _ = tx_processed.send(result).await;
                                                }
                                                None => break,
                                            }
                                        }
                                        _ = shutdown_rx.recv() => break,
                                    }
                                }
                                break;
                            }
                        }
                    }
                    // 响应关闭信号
                    _ = shutdown_rx.recv() => break,
                }
            }
        });

        Self {
            tx_to_alice,
            tx_to_bob,
            rx_processed,
            shutdown_tx,
        }
    }

    async fn send(&self, message: Message, dest: Destination) -> Result<()> {
        match dest {
            Destination::Alice => self.tx_to_alice.send(message)
                .map_err(|e| anyhow!("发送消息到Alice失败: {}", e)),
            Destination::Bob => self.tx_to_bob.send(message)
                .map_err(|e| anyhow!("发送消息到Bob失败: {}", e)),
        }
    }

    // 获取下一条处理完成的消息
    async fn next_processed(&mut self) -> Result<Message> {
        self.rx_processed.recv().await
            .ok_or_else(|| anyhow!("无更多处理结果,所有连接已关闭"))?
    }

    // 异步消息处理逻辑
    async fn process_received_message(message: Message) -> Result<Message> {
        println!("正在处理消息: {:?}", message);
        Ok(message)
    }

    // 检查Alice连接状态
    pub fn alice_is_closed(&self) -> bool {
        self.tx_to_alice.is_closed()
    }

    // 检查Bob连接状态
    pub fn bob_is_closed(&self) -> bool {
        self.tx_to_bob.is_closed()
    }
}

// 使用示例
async fn listen_to_a_lot_of_stuff() -> Result<()> {
    let mut connections_manager = ConnectionsManager::new();

    loop {
        tokio::select! {
            processed_message = connections_manager.next_processed() => {
                match processed_message {
                    Ok(msg) => println!("收到处理完成的消息: {:?}", msg),
                    Err(e) => {
                        println!("消息处理出错: {}", e);
                        break;
                    }
                }
            }
            // 其他业务逻辑示例
            _ = tokio::time::sleep(tokio::time::Duration::from_secs(5)) => {
                println!("执行其他任务...");
                // 可正常调用send方法,无引用冲突
                connections_manager.send(Message("来自主逻辑的Ping".into()), Destination::Alice).await?;
            }
        }
    }

    Ok(())
}

方案优势

  • 完全避免消息丢失:后台任务持续运行,只要消息从Alice/Bob的通道被接收,就一定会被处理,不受外层select!切换影响。
  • 无引用冲突:send()仅需要不可变引用,与后台任务运行互不干扰,不存在持有Future导致的可变引用锁定问题。
  • 可追踪连接状态:通过tx_to_alice.is_closed()/tx_to_bob.is_closed()可直接判断对应连接是否关闭,无需依赖RX通道的None结果。

补充说明

  • 内部处理结果通道使用有界通道,可防止处理速度跟不上接收速度导致的内存溢出问题,容量可根据业务调整。
  • 后台任务的关闭信号可确保ConnectionsManager销毁时优雅停止任务,避免资源泄漏。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 06:25:23