如何在futures::stream::Stream并发处理时终止中间Future?
问题:多Peer并发任务中终止所有task#2的实现
我尝试并发连接多个Peer并进行处理,单个任务(称为“task”)运行正常。但在多任务场景下,我希望当某个Peer完成task#2后,终止所有其他Peer的task#2,包括来自create_new_conn_fut的后续Peer。
现有代码
use futures::stream::StreamExt; use rand::Rng; pub async fn process(peer: &str, duration: core::time::Duration, task_id: &str) { // simulate processing by sleeping tokio::time::sleep(duration).await; println!("task #{} done for {}", task_id, peer); } #[tokio::main] async fn main() { let peers = vec!["peer A", "peer B", "peer C"]; let peers = futures::stream::iter(peers); let (tx, rx) = tokio::sync::mpsc::channel(100); let rx = tokio_stream::wrappers::ReceiverStream::new(rx); let rx = peers.chain(rx); let handle_conn_fut = rx.for_each_concurrent(0, |peer| async move { let mut rng = rand::thread_rng(); println!("connecting to {}", peer); process(peer, core::time::Duration::from_secs(1), "1").await; process(peer, core::time::Duration::from_secs(rng.gen_range(5..15)), "2").await; process(peer, core::time::Duration::from_secs(1), "3").await; } ); let create_new_conn_fut = async move { for peer in ["peer D", "peer E"] { tx.send(peer).await.unwrap(); } }; // awaits all futures in parallell futures::future::join(handle_conn_fut, create_new_conn_fut).await; }
当前输出
connecting to peer A connecting to peer B connecting to peer C connecting to peer D connecting to peer E task #1 done for peer A task #1 done for peer B task #1 done for peer C task #1 done for peer D task #1 done for peer E task #2 done for peer C task #3 done for peer C task #2 done for peer D task #2 done for peer A task #2 done for peer B task #3 done for peer D task #3 done for peer A task #3 done for peer B task #2 done for peer E task #3 done for peer E
期望输出
connecting to peer A connecting to peer B connecting to peer C connecting to peer D connecting to peer E task #1 done for peer A task #1 done for peer B task #1 done for peer C task #1 done for peer D task #1 done for peer E task #2 done for peer C <- will abort all other task #2 task #3 done for peer A task #3 done for peer B task #3 done for peer C task #3 done for peer D task #3 done for peer E
我曾研究过futures::future::AbortHandle,但认为它仅适用于单个Future,因为futures::stream::AbortRegistration不具备Clone trait。请问该如何实现此需求?
解决方案
可以使用**tokio::sync::Notify**实现全局通知:当第一个Peer完成task#2后,通过Notify广播信号,所有正在执行或即将执行task#2的Peer收到信号后直接跳过task#2,进入task#3。
修改后的代码
use futures::stream::StreamExt; use rand::Rng; use tokio::sync::Notify; pub async fn process(peer: &str, duration: core::time::Duration, task_id: &str) { tokio::time::sleep(duration).await; println!("task #{} done for {}", task_id, peer); } #[tokio::main] async fn main() { let peers = vec!["peer A", "peer B", "peer C"]; let peers = futures::stream::iter(peers); let (tx, rx) = tokio::sync::mpsc::channel(100); let rx = tokio_stream::wrappers::ReceiverStream::new(rx); let rx = peers.chain(rx); // 创建全局通知器,用于广播task#2完成的信号 let task2_done_notify = Notify::new(); // 克隆通知器,供每个Peer任务使用 let task2_done_clone = task2_done_notify.clone(); let handle_conn_fut = rx.for_each_concurrent(0, move |peer| { let notify = task2_done_clone.clone(); let done_notify = task2_done_notify.clone(); async move { let mut rng = rand::thread_rng(); println!("connecting to {}", peer); process(peer, core::time::Duration::from_secs(1), "1").await; // 检查是否已有Peer完成task#2,若没有则执行,执行完后通知所有Peer if !notify.is_notified() { let task2_fut = process(peer, core::time::Duration::from_secs(rng.gen_range(5..15)), "2"); tokio::select! { _ = task2_fut => { println!("task #2 done for {}", peer); // 通知所有等待的Peer,task#2已完成 done_notify.notify_waiters(); } _ = notified(¬ify) => { // 收到通知,跳过task#2 println!("skipping task #2 for {}", peer); } } } else { // 已有Peer完成task#2,直接跳过 println!("skipping task #2 for {}", peer); } process(peer, core::time::Duration::from_secs(1), "3").await; } }); let create_new_conn_fut = async move { for peer in ["peer D", "peer E"] { tx.send(peer).await.unwrap(); } }; futures::future::join(handle_conn_fut, create_new_conn_fut).await; } // 辅助函数,等待Notify通知 async fn notified(notify: &Notify) { notify.notified().await; }
实现说明
- 全局通知器
Notify:tokio::sync::Notify是多生产者多消费者的通知机制,支持广播信号和Clone,适配多任务场景。 tokio::select!分支:执行task#2时同时监听两个事件:- task#2自身完成:触发全局通知,告知所有Peer跳过task#2。
- 收到全局通知:直接放弃执行task#2,进入task#3。
- 前置检查:后续加入的Peer先通过
is_notified()判断是否已有Peer完成task#2,若是则直接跳过task#2。
这样就能实现需求:第一个完成task#2的Peer会终止所有其他Peer的task#2,后续新Peer也不会执行task#2。
内容的提问来源于stack exchange,提问作者Kevin
相关产品推荐
相关产品推荐

