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

Rust中如何同时监听多组异步流与停止消息?

同时监听多组异步接收器与停止信号的解决方案

核心思路

要同时监听所有接收器对和停止信号,需要将每组(Receiver<SomeType>, Receiver<SomeOtherType>)的消息合并为统一的流,再将所有流合并成一个全局流,最后与停止信号流一起监听。

具体实现(基于Tokio)

首先定义枚举统一两种消息类型,方便后续处理:

enum BusMessage {
    Sample(SomeType),
    Event(SomeOtherType),
    Stop,
}

然后将每组接收器转换为流,合并其中的样本和事件消息:

use tokio_stream::{Stream, StreamExt, wrappers::ReceiverStream};
use futures::stream::select;

fn combine_pair(rx_sample: Receiver<SomeType>, rx_event: Receiver<SomeOtherType>) -> impl Stream<Item = BusMessage> {
    let sample_stream = ReceiverStream::new(rx_sample).map(BusMessage::Sample);
    let event_stream = ReceiverStream::new(rx_event).map(BusMessage::Event);
    select(sample_stream, event_stream)
}

接下来合并所有流,并与停止信号一起监听:

use futures::stream::select_all;
use tokio::select;

async fn run_receivers(
    receivers: Vec<(Receiver<SomeType>, Receiver<SomeOtherType>)>,
    stop_rx: Receiver<()>
) {
    // 合并所有接收器对的流
    let mut combined_streams = select_all(
        receivers.into_iter().map(|pair| combine_pair(pair.0, pair.1))
    );

    // 将停止信号转换为流
    let mut stop_stream = ReceiverStream::new(stop_rx).map(|_| BusMessage::Stop);

    loop {
        match select! {
            msg = combined_streams.next() => {
                match msg {
                    Some(BusMessage::Sample(sample)) => {
                        // 处理数据样本
                        println!("Got sample: {:?}", sample);
                    }
                    Some(BusMessage::Event(event)) => {
                        // 处理相关事件
                        println!("Got event: {:?}", event);
                    }
                    None => {
                        // 所有接收器都已关闭,退出循环
                        break;
                    }
                }
            }
            Some(BusMessage::Stop) = stop_stream.next() => {
                // 收到停止信号,终止接收器
                println!("Stopping receiver...");
                break;
            }
            None => break,
        }
    }
}

关键细节说明

  1. 为什么用流而不是直接处理recv() Future:
    Tokio的Receiver的recv()方法需要&mut self,无法同时生成多个独立的recv() Future。将接收器转换为ReceiverStream后,流会自动处理连续的消息接收,避免了所有权和可变引用的问题。

  2. 解决FusedFuture绑定错误:
    select_all要求输入的流实现FusedStream,而ReceiverStream天然满足这个 trait。之前的错误通常是因为直接将未包装的recv() Future传入select_all,这类Future不满足FusedFuture绑定。

  3. 自动处理关闭的接收器:
    当某个接收器关闭时,select_all会自动将其从合并流中移除,无需手动管理,剩下的接收器会继续被监听。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 17:40:26