Rust实现Java CompletableFuture::complete功能及mpsc使用合理性咨询
Rust实现类似Java CompletableFuture::complete的功能
Rust完全可以实现你需要的「在外部线程完成Future」的需求,通道是符合Rust异步设计理念的常用方案,针对你的场景,我们可以从以下几个方面梳理:
为什么oneshot版本无法运行
oneshot::Sender的设计是一次性、需要所有权才能发送,它既不支持Clone,也不能通过共享引用调用send方法。你尝试用RwLock包裹Sender并共享,但send会转移Sender的所有权,导致锁内的Sender被移动,这违背了RwLock的使用逻辑(不能在持有读锁时修改内部数据)。这是oneshot通道的特性,它只允许唯一的Sender发送一次数据,所以这种共享存储Sender的方式不适合oneshot。
mpsc版本的合理性与性能
你用容量为1的mpsc通道实现的方案是合理且高效的:
- 容量1完全匹配你「只发送一次结果」的需求,mpsc通道在低容量下的性能损耗极小,和oneshot通道几乎无差异;
- mpsc的
Sender支持Clone,这让你可以在Executor中存储Sender的克隆,后续通过索引获取后发送数据; - 你代码中的手动
close其实可以省略:当所有Sender实例被drop后,Receiver的recv方法会返回None,所以不需要手动调用remove,run方法里的循环也可以简化为直接等待一次接收。
优化后的mpsc版本代码:
#[tokio::test] pub async fn test_async() -> Result<(), Box<dyn std::error::Error>> { let mut executor = Executor::new(); let res1 = executor.run(1).await?; let res2 = executor.run(2).await?; println!("res1 {}, res2 {}", res1, res2); Ok(()) } struct Executor { pub pool: rayon::ThreadPool, pub txs: Arc<dashmap::DashMap<i32, mpsc::Sender<i32>>>, } impl Executor { pub fn new() -> Self { Executor { pool: rayon::ThreadPoolBuilder::new().num_threads(10).build().unwrap(), txs: Arc::new(dashmap::DashMap::new()), } } pub async fn run(&mut self, index: i32) -> Result<i32, Box<dyn std::error::Error>> { let (tx, mut rx) = mpsc::channel(1); let tx_clone = tx.clone(); self.txs.insert(index, tx_clone); let txs_clone = self.txs.clone(); self.pool.spawn(move || { if let Some(tx) = txs_clone.get(&index) { tokio::runtime::Handle::current().block_on(tx.send(index)).unwrap(); } }); // 直接等待一次接收,因为只发送一次 let res = rx.recv().await.ok_or("未收到结果")?; // 接收完成后移除Sender,避免内存泄漏 self.txs.remove(&index); Ok(res) } }
更贴合场景的oneshot优化方案
如果你更倾向于用oneshot(毕竟是一次性发送的场景),可以通过Arc<Mutex<Option<oneshot::Sender<T>>>>来存储Sender,在需要发送时将Option中的Sender取出(获得所有权)再发送:
#[tokio::test] pub async fn test_async() -> Result<(), Box<dyn std::error::Error>> { let mut executor = Executor::new(); let res1 = executor.run(1).await?; let res2 = executor.run(2).await?; println!("res1 {}, res2 {}", res1, res2); Ok(()) } struct Executor { pub pool: rayon::ThreadPool, pub txs: Arc<dashmap::DashMap<i32, Arc<tokio::sync::Mutex<Option<oneshot::Sender<i32>>>>>>, } impl Executor { pub fn new() -> Self { Executor { pool: rayon::ThreadPoolBuilder::new().num_threads(10).build().unwrap(), txs: Arc::new(dashmap::DashMap::new()), } } pub async fn run(&mut self, index: i32) -> Result<i32, Box<dyn std::error::Error>> { let (tx, rx) = oneshot::channel(); self.txs.insert(index, Arc::new(tokio::sync::Mutex::new(Some(tx)))); let txs_clone = self.txs.clone(); self.pool.spawn(move || { if let Some(entry) = txs_clone.get(&index) { let mut guard = tokio::runtime::Handle::current().block_on(entry.lock()); // 取出Sender获得所有权 if let Some(sender) = guard.take() { sender.send(index).unwrap(); } } }); let res = rx.await?; self.txs.remove(&index); Ok(res) } }
总结
- 容量为1的mpsc方案完全可行,性能足够,适合你的场景;
- 优化后的oneshot方案也能满足需求,更贴合「一次性发送结果」的语义;
- 两种方案都不需要手动管理通道关闭,在接收完成或Sender被取出后清理存储即可避免内存泄漏。
内容的提问来源于stack exchange,提问作者lulijun
相关产品推荐
相关产品推荐

