如何轮询Pin<Option<Box<dyn Future>>>?异步递归流实现问题
问题:如何正确轮询
Pin<Option<Box<dyn Future>>>以实现递归异步函数生成的Stream 相关代码如下:
use std::{ future::Future, pin::Pin, task::{Context, Poll}, }; use futures_util::Stream; #[pin_project::pin_project] struct BoxedStream<T, F> { #[pin] next: Option<String>, #[pin] inner: Option<Box<dyn Future<Output = T>>>, // ^--- 异步函数生成的Future,是不透明类型 generate: F, } impl<T, F> Stream for BoxedStream<T, F> where F: Fn(String) -> (Option<String>, Box<dyn Future<Output = T>>), { type Item = T; fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> { let mut p = self.as_mut().project(); if let Some(s) = p.inner.as_mut().as_pin_mut() { return match todo!("how to do it here? s.poll(cx)?") { result @ Poll::Ready(_) => { p.inner.take(); result } _ => Poll::Pending, }; } if let Some(next) = p.next.take() { let (next, future) = (p.generate)(next); p.inner.set(Some(future)); p.next.set(next); return self.poll_next(cx); } Poll::Ready(None) // 结束 } }
调用s.poll(cx)替代todo代码时,出现如下E0599错误:
error[E0599]: the method `poll` exists for struct `Pin<&mut Box<dyn Future<Output = T>>>`, but its trait bounds were not satisfied --> src/lib.rs:28:28 | 28 | return match s.poll(cx)? { | ^^^^ method cannot be called on `Pin<&mut Box<dyn Future<Output = T>>>` due to unsatisfied trait bounds --> /rustc/5680fa18feaa87f3ff04063800aec256c3d4b4be/library/core/src/future/future.rs:37:1 | = note: doesn't satisfy `_: Unpin` --> /rustc/5680fa18feaa87f3ff04063800aec256c3d4b4be/library/alloc/src/boxed.rs:195:1 ::: /rustc/5680fa18feaa87f3ff04063800aec256c3d4b4be/library/alloc/src/boxed.rs:198:1 | = note: doesn't satisfy `_: Future` | = note: the following trait bounds were not satisfied: `(dyn futures_util::Future<Output = T> + 'static): Unpin` which is required by `Box<(dyn futures_util::Future<Output = T> + 'static)>: futures_util::Future`
错误原因分析
错误核心是Pin<&mut Box<dyn Future<Output = T>>>不满足Future trait的约束:
Box<dyn Future>要实现Future,要求内部的dyn Future必须满足Unpin,但异步函数生成的匿名Future默认不自动实现Unpin- 当前代码的
Pin包装的是Box<dyn Future>,但我们需要直接把Pin作用到内部的dyn Future上,才能调用poll方法
解决方案
方案1:给dyn Future添加Unpin约束(简单直接)
如果异步函数生成的Future可以实现Unpin(比如使用Box::pin或允许Unpin),直接修改类型定义:
#[pin_project::pin_project] struct BoxedStream<T, F> { #[pin] next: Option<String>, #[pin] inner: Option<Box<dyn Future<Output = T> + Unpin>>, // 添加Unpin约束 generate: F, } impl<T, F> Stream for BoxedStream<T, F> where // 同步修改闭包返回类型的约束 F: Fn(String) -> (Option<String>, Box<dyn Future<Output = T> + Unpin>), { type Item = T; fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> { let mut p = self.as_mut().project(); if let Some(s) = p.inner.as_mut().as_pin_mut() { return match s.poll(cx) { Poll::Ready(item) => { p.inner.take(); Poll::Ready(Some(item)) } Poll::Pending => Poll::Pending, }; } if let Some(next) = p.next.take() { let (next, future) = (p.generate)(next); p.inner.set(Some(future)); p.next.set(next); return self.poll_next(cx); } Poll::Ready(None) } }
方案2:手动穿透Pin到内部dyn Future(通用无约束方案)
如果无法添加Unpin约束,手动将Pin穿透到Box内部的dyn Future上:
fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> { let mut p = self.as_mut().project(); if let Some(boxed_fut) = p.inner.as_mut().as_pin_mut() { // 将Pin<&mut Box<dyn Future>>转换为Pin<&mut dyn Future> let pinned_fut = Pin::as_mut(boxed_fut).as_mut(); return match pinned_fut.poll(cx) { Poll::Ready(item) => { p.inner.take(); Poll::Ready(Some(item)) } Poll::Pending => Poll::Pending, }; } if let Some(next) = p.next.take() { let (next, future) = (p.generate)(next); p.inner.set(Some(future)); p.next.set(next); return self.poll_next(cx); } Poll::Ready(None) }
这里利用Pin::as_mut()将外层的Pin<&mut Box<F>>转换为直接作用于内部Future的Pin<&mut F>,无需Unpin约束即可调用poll。
方案3:使用BoxFuture简化实现(推荐)
futures_util提供的BoxFuture是Pin<Box<dyn Future<Output = T> + Send + 'static>>的别名,直接支持poll操作,无需手动处理Pin逻辑:
use futures_util::{future::BoxFuture, Stream}; #[pin_project::pin_project] struct BoxedStream<T, F> { #[pin] next: Option<String>, #[pin] inner: Option<BoxFuture<'static, T>>, // 使用BoxFuture替代手动Box<dyn Future> generate: F, } impl<T, F> Stream for BoxedStream<T, F> where F: Fn(String) -> (Option<String>, BoxFuture<'static, T>), { type Item = T; fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> { let mut p = self.as_mut().project(); if let Some(fut) = p.inner.as_mut().as_pin_mut() { return match fut.poll(cx) { Poll::Ready(item) => { p.inner.take(); Poll::Ready(Some(item)) } Poll::Pending => Poll::Pending, }; } if let Some(next) = p.next.take() { let (next, future) = (p.generate)(next); p.inner.set(Some(future)); p.next.set(next); return self.poll_next(cx); } Poll::Ready(None) } }
生成Future时,用Box::pin(async { ... })创建BoxFuture:
fn generate_future(s: String) -> (Option<String>, BoxFuture<'static, i32>) { (Some("next".to_string()), Box::pin(async { // 异步逻辑 s.len() as i32 })) }
内容的提问来源于stack exchange,提问作者holi-java
相关产品推荐
相关产品推荐

