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
相关产品推荐
相关产品推荐

