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

如何同步接收不同类型async_channel消息以操作共享内存?

如何同时从不同类型的async_channel接收消息?

我希望通过同步访问共享内存shared简化代码,该内存会根据异步到达的请求/事件修改。elsewhere::run函数等待另一线程中sender发来的请求,请求为不同类型结构体(Foo和Bar),通过async_channel发送,每种请求对共享内存有不同操作。我尝试用futures::future::select_all绑定不同的async_channel::Receiver::recv()调用,但无法实现。需要同时等待两个接收器,收到第一条消息时同步恢复任务,请问该如何实现?

以下是最小可复现示例:

use async_channel::Sender;
use std::{thread, time};
use futures::executor::block_on;

pub struct Foo {}
pub struct Bar {}

pub mod elsewhere {
    use futures::future::select_all;
    use crate::{Foo, Bar};
    use async_channel::Receiver;
    
    enum FooBar {
        Foo(Foo),
        Bar(Bar)
    }
    
    pub async fn run(rx_foo: Receiver<Foo>, rx_bar: Receiver<Bar>) {
        println!("run is listening");
        
        let mut shared: Vec<i32> = vec![0]; // 尝试通过同步访问避免给共享内存加互斥锁
        
        loop {
            let _recv_foo = rx_foo.recv();
            let _recv_bar = rx_bar.recv();
            
            /**
             * [...]
             * 将Foo和Bar转换为FooBar,并用类似"select_all"的方式等待
             * 两个通道rx_foo & rx_bar中最先到达的消息
             */
            
            let firstReceived = FooBar::Foo(Foo{}); // 用 dummy 值替代通道消息以保证编译通过
            
            match firstReceived {
                Foo => { shared.push(0);}
                Bar => { shared.clear();}
            }
        }
    }
}

async fn sender(tx_foo: Sender<Foo>, tx_bar: Sender<Bar>) {
    println!("Send to run");
    let _ = tx_foo.send(Foo{}).await;
    let _ = tx_bar.send(Bar{}).await;
}

fn main() {        
    // 希望main模块无需知晓elsewhere::FooBar
    let (tx_foo, rx_foo) = async_channel::unbounded::<Foo>();
    let (tx_bar, rx_bar) = async_channel::unbounded::<Bar>();

    thread::spawn(move || block_on(elsewhere::run(rx_foo, rx_bar)));
    thread::spawn(move || block_on(sender(tx_foo, tx_bar)));

    thread::sleep(time::Duration::from_secs(10));
}

解决方案

核心思路是把两个不同类型的recv() Future转换成返回相同枚举类型的Future,这样就可以用futures::future::select(针对两个Future)或select_all(针对多个Future)来等待最先完成的那个。

方法一:使用select(适合两个通道场景)

修改后的elsewhere::run函数:

pub mod elsewhere {
    use futures::future::{select, Either};
    use crate::{Foo, Bar};
    use async_channel::Receiver;
    
    enum FooBar {
        Foo(Foo),
        Bar(Bar)
    }
    
    pub async fn run(mut rx_foo: Receiver<Foo>, mut rx_bar: Receiver<Bar>) {
        println!("run is listening");
        
        let mut shared: Vec<i32> = vec![0];
        
        loop {
            // 将两个recv Future映射为返回FooBar的Future,统一输出类型
            let recv_foo = rx_foo.recv().map(FooBar::Foo);
            let recv_bar = rx_bar.recv().map(FooBar::Bar);
            
            // 用select等待最先到达的消息
            match select(recv_foo, recv_bar).await {
                // 处理foo通道先收到消息的情况
                Either::Left((Ok(foo_bar), _unused_future)) => {
                    match foo_bar {
                        FooBar::Foo(_) => shared.push(0),
                        FooBar::Bar(_) => shared.clear(),
                    }
                }
                // 处理bar通道先收到消息的情况
                Either::Right((Ok(foo_bar), _unused_future)) => {
                    match foo_bar {
                        FooBar::Foo(_) => shared.push(0),
                        FooBar::Bar(_) => shared.clear(),
                    }
                }
                // 可选:处理通道关闭的异常情况
                Either::Left((Err(_), _)) | Either::Right((Err(_), _)) => {
                    println!("其中一个通道已关闭");
                    break;
                }
            }
        }
    }
}

方法二:使用select_all(适合多通道场景)

如果后续需要扩展更多通道,用select_all更灵活:

pub mod elsewhere {
    use futures::future::select_all;
    use futures::{Future, pin_mut};
    use std::pin::Pin;
    use crate::{Foo, Bar};
    use async_channel::Receiver;
    
    enum FooBar {
        Foo(Foo),
        Bar(Bar)
    }
    
    pub async fn run(mut rx_foo: Receiver<Foo>, mut rx_bar: Receiver<Bar>) {
        println!("run is listening");
        
        let mut shared: Vec<i32> = vec![0];
        
        loop {
            // 把多个统一类型的Future放入集合
            let futures: Vec<Pin<Box<dyn Future<Output = Result<FooBar, async_channel::RecvError>>>>> = vec![
                Box::pin(rx_foo.recv().map(FooBar::Foo)),
                Box::pin(rx_bar.recv().map(FooBar::Bar)),
            ];
            
            match select_all(futures).await {
                (Ok(foo_bar), _index, _remaining) => {
                    match foo_bar {
                        FooBar::Foo(_) => shared.push(0),
                        FooBar::Bar(_) => shared.clear(),
                    }
                }
                (Err(_), _index, _remaining) => {
                    println!("某个通道已关闭");
                    break;
                }
            }
        }
    }
}

关键说明

  • 用.map()把recv()返回的Result<Foo, ...>和Result<Bar, ...>转换成Result<FooBar, ...>,让两个Future的输出类型统一,满足select/select_all的类型要求。
  • 共享内存shared可以安全同步访问,因为所有修改逻辑都在同一个异步任务的上下文里执行,无需互斥锁。
  • 保持了main模块对FooBar的无知,类型转换逻辑完全封装在elsewhere模块内部,符合代码隔离需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 06:04:55