如何同步接收不同类型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
相关产品推荐
相关产品推荐

