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

Rust异步回调加入FuturesUnordered队列的类型检查报错问题

问题说明

需要实现一个接收异步回调、将回调生成的异步任务放入FuturesUnordered队列的函数。现有代码运行逻辑完全符合预期,但无法通过类型检查,初步猜测问题与Pin、Box、dyn相关,需要确认该需求是否可以实现。

已尝试过多种修改方案,始终无法编译通过,不希望通过盲目试错修改代码,核心诉求为确认该逻辑的实现可行性:

  • 调整生命周期标注
  • 添加不同形式的Box与Pin包装
  • 遵循编译器提示修改代码

初步判断Rust无法验证代码中引用的内存安全性,需要将部分对象进行装箱操作,为编译器提供固定(pinned)的常量访问权限,但无法确定需要对哪些对象做包装处理。


原始实现代码

use async_stream::stream;
use futures::stream::{FuturesUnordered, StreamExt};

pub fn map<U, V, W>(
    f: impl Fn(&U) -> W,
    items: Vec<U>,
) -> impl futures::Stream<Item = V>
where
    V: Send,
    W: futures::Future<Output = V> + Send,
{
    stream! {
        let mut futures = FuturesUnordered::new();
        let mut i = 2;
        if 2 <= items.len() {
            futures.push(tokio::spawn(f(&items[0])));
            futures.push(tokio::spawn(f(&items[1])));
            while let Some(result) = futures.next().await {
                let y = result.unwrap();
                yield y;
                futures.push(tokio::spawn(f(&items[i])));
                i += 1
            }
        }
    }
}

#[tokio::main]
async fn main() {
    async fn f(x: &u32) -> u32 {
        x + 1
    }
    let input = vec![1, 2, 3];
    let output = map(f, input);
    futures::pin_mut!(output);
    while let Some(x) = output.next().await {
        println!("{:?}", x);
    }
}

原始编译报错

编译上述代码会报E0308类型不匹配错误,具体错误信息如下:

error[E0308]: mismatched types
  --> src/main.rs:40:18
   |
40 |     let output = map(f, input);
   |                  ^^^ lifetime mismatch
   |
   = note: expected associated type `<for<'_> fn(&u32) -> impl futures::Future<Output = u32> {f} as FnOnce<(&u32,)>>::Output`
              found associated type `<for<'_> fn(&u32) -> impl futures::Future<Output = u32> {f} as FnOnce<(&u32,)>>::Output`
   = note: the required lifetime does not necessarily outlive the empty lifetime

已尝试修改方案与对应报错

其中一个修改方案为将map函数的签名修改为回调返回Pin<Box<dyn Future<Output = V> + '_>>的形式,修改后的函数签名如下:

pub fn map<U, V>(
    f: impl Fn(&U) -> Pin<Box<dyn Future<Output = V> + '_>>,
    items: Vec<U>,
) -> impl futures::Stream<Item = V>
where
    V: Send,

修改后Rust报E0277错误,提示tokio::spawn要求的Send trait未实现,具体错误信息如下:

error[E0277]: `dyn futures::Future<Output = V>` cannot be sent between threads safely
   --> src/main.rs:21:26
    |
21  |             futures.push(tokio::spawn(f(&items[0])));
    |                          ^^^^^^^^^^^^ `dyn futures::Future<Output = V>` cannot be sent between threads safely
    |
    = help: the trait `std::marker::Send` is not implemented for `dyn futures::Future<Output = V>`
    = note: required because of the requirements on the impl of `std::marker::Send` for `Unique<dyn futures::Future<Output = V>>`
    = note: required because it appears within the type `Box<dyn futures::Future<Output = V>>`
    = note: required because it appears within the type `Pin<Box<dyn futures::Future<Output = V>>>`

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 07:09:15