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

如何在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(&notify) => {
                        // 收到通知,跳过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;
}

实现说明

  1. 全局通知器Notify:tokio::sync::Notify是多生产者多消费者的通知机制,支持广播信号和Clone,适配多任务场景。
  2. tokio::select!分支:执行task#2时同时监听两个事件:
    • task#2自身完成:触发全局通知,告知所有Peer跳过task#2。
    • 收到全局通知:直接放弃执行task#2,进入task#3。
  3. 前置检查:后续加入的Peer先通过is_notified()判断是否已有Peer完成task#2,若是则直接跳过task#2。

这样就能实现需求:第一个完成task#2的Peer会终止所有其他Peer的task#2,后续新Peer也不会执行task#2。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 22:26:25