在编译时等待许多未知的 future

arm*_*ani 5 rust async-await rust-tokio

我想利用 Tokio 的运行时来处理可变数量的异步 future。由于 futures 的数量在编译时是未知的,似乎FuturesUnordered是我最好的选择(宏,例如select!需要在编译时指定您的分支;join_all可能是可能的,但文档建议“在很多情况下”当 order 不可用时使用 FuturesUnordered )没关系)。

这段代码的逻辑是一个recv()循环被推送到futures桶中,它应该始终运行。当新数据到达时,它的解析/处理也被推送到 futures 桶(而不是立即处理)。这确保接收器在响应新事件时保持低延迟,并且数据处理(可能需要大量计算的解密)与所有其他数据处理异步块(加上侦听接收器)同时发生。

.boxed()顺便说一句,这个帖子解释了为什么期货会变得。

问题是这个神秘的错误:

错误[E0277] :`dyn futures::Future<Output = ()> + std::marker::Send` 无法在线程之间安全共享
  --> src/main.rs:27:8
    | 
27  | 27     }).boxed());
   |        ^^^^^  `dyn futures::Future<Output = ()> + std::marker::Send` 无法在线程之间安全共享
   | 
   = help : `dyn futures::Future<Output = ()> + std::marker::Send` 未实现 `Sync` 特性
   =注意:需要,因为 `Sync` 的 impl 要求Unique<dyn futures::Future<Output = ()> + std::marker::Send>`
    = note:必需的,因为它出现在类型 `Box<dyn futures::Future<Output = ()> + std: :marker::Send>`
    = note : 必需,因为它出现在类型 `Pin<Box<dyn futures::Future<Output = ()> + std::marker::Send>>`
    = note : 必需,因为`FuturesUnordered<Pin<Box<dyn futures::Future<Output = ()> + std::marker::Send>>>` 的 `Sync` 实现的要求
   =注意:需要,因为对`&FuturesUnordered<Pin<Box<dyn futures::Future<Output = ()> + std::marker::Send>>>` 的 `std::marker::Send` 的实现 =
   注意:必需,因为它出现在类型 `[static Generator@src/main.rs:16:25: 27:6 _]`
    = note : 必需,因为它出现在类型 `from_generator::GenFuture<[static Generator@src/main.rs:16 :25: 27:6 _]>`
    =注意:必需,因为它出现在 `impl futures::Future` 类型中

看起来“递归地”推送到 UnorderedFutures (我猜不是真的,但你还能称之为什么?)不起作用,但我不确定为什么。此错误表明SyncBox'd 和 Pin'd 异步块未满足某些特征要求FuturesUnordered- 我猜这个要求只是强加的,因为&FuturesUnordered(在该方法借用 &self 期间使用futures.push(...))需要它的Send特征... 或者其他的东西?

use std::error::Error;
use tokio::sync::mpsc::{self, Receiver, Sender};
use futures::stream::futures_unordered::FuturesUnordered;
use futures::FutureExt;

#[tokio::main]
pub async fn main() -> Result<(), Box<dyn Error>> {
    let mut futures = FuturesUnordered::new();
    let (tx, rx) = mpsc::channel(32);
    
    tokio::spawn( foo(tx) );    // Only the receiver is relevant; its transmitter is
                                // elsewhere, occasionally sending data.
    futures.push((async {                               // <--- NOTE: futures.push()
        loop {
            match rx.recv().await {
                Some(data) => {
                    futures.push((async move {          // <--- NOTE: nested futures.push()
                        let _ = data; // TODO: replace with code that processes 'data'
                    }).boxed());
                },
                None => {}
            }
        }
    }).boxed());
    
    while let Some(_) = futures.next().await {}

    Ok(())
}
Run Code Online (Sandbox Code Playgroud)

all*_*y87 10

我将把低级错误留给另一个答案,但我相信解决高级问题的更惯用的方法是将FuturesUnordered与 的使用结合起来tokio::select!,如下所示:

use tokio::sync::mpsc;
use futures::stream::FuturesUnordered;
use futures::StreamExt;

#[tokio::main]
pub async fn main() {
    let mut futures = FuturesUnordered::new();
    let (tx, mut rx) = mpsc::channel(32);
    
    //turn foo into something more concrete
    tokio::spawn(async move {
        let _ = tx.send(42i32).await;
    });

    loop {
        tokio::select! {
            Some(data) = rx.recv() => {
                futures.push(async move {
                    data.to_string()
                });
            },
            Some(result) = futures.next() => {
                println!("{}", result)
            },
            else => break,
        }
    }
}
Run Code Online (Sandbox Code Playgroud)

您可以在此处阅读有关 select 宏的更多信息: https: //tokio.rs/tokio/tutorial/select