如何将与API源生命周期绑定的Stream转换为'static Stream?
如何将绑定到Arc的Stream转换为'static生命周期
问题背景
给定如下API:
struct API; struct Ctrl; impl API{ fn get_stream(&self) -> (impl Stream<Item = i32> + '_, Ctrl) { (futures::stream::iter(1..=5), Ctrl{}) // 占位流 } }
该API返回的Stream生命周期与&self绑定,而我们需要一个'static的Stream以便移出当前作用域。尝试通过封装Combine结构体持有Arc<API>和Stream的方案失败,因为编译器无法识别Arc能保证Stream依赖的API实例生命周期——调用instance.get_stream()时,&instance是临时引用,Stream的生命周期绑定到这个临时引用,而非Arc内部的API实例。
解决方案
方案一:使用async_stream crate创建'static Stream
通过async_stream创建新Stream,捕获Arc<API>实例,内部消费原Stream并yield元素。新Stream因持有Arc<API>,天然满足'static生命周期要求。
代码示例:
use async_stream::stream; use futures::stream::Stream; use futures::StreamExt; use std::sync::Arc; struct API; struct Ctrl; impl API { fn get_stream(&self) -> (impl Stream<Item = i32> + '_ , Ctrl) { (futures::stream::iter(1..=5), Ctrl{}) } } fn wrap_static_stream(instance: Arc<API>) -> impl Stream<Item = i32> + 'static { stream! { let (mut inner_stream, _ctrl) = instance.get_stream(); while let Some(item) = inner_stream.next().await { yield item; } } } #[tokio::main] async fn main() { let static_stream = { let instance = Arc::new(API{}); wrap_static_stream(instance) }; tokio::pin!(static_stream); while let Some(item) = static_stream.next().await { println!("Received item: {}", item); } }
方案二:用ouroboros定义自引用结构体
借助ouroboros库生成合法的自引用结构体,明确标注Stream与Arc<API>的生命周期绑定,解决编译器推导问题。
首先在Cargo.toml添加依赖:
ouroboros = { version = "0.15", features = ["futures"] }
代码示例:
use futures::stream::{Stream, StreamExt}; use futures::task::{Context, Poll}; use ouroboros::self_referencing; use std::pin::Pin; use std::sync::Arc; struct API; struct Ctrl; impl API { fn get_stream(&self) -> (impl Stream<Item = i32> + '_ , Ctrl) { (futures::stream::iter(1..=5), Ctrl{}) } } #[self_referencing] struct StaticStream { instance: Arc<API>, #[borrows(instance)] #[covariant] stream: Pin<Box<dyn Stream<Item = i32> + 'this>>, } impl Stream for StaticStream { type Item = i32; fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> { self.with_mut(|fields| { Pin::as_mut(&mut fields.stream).poll_next(cx) }) } } fn make_static_stream(instance: Arc<API>) -> StaticStream { StaticStreamBuilder { instance, stream_builder: |instance_ref| { let (stream, _ctrl) = instance_ref.get_stream(); Box::pin(stream) }, }.build() } #[tokio::main] async fn main() { let static_stream = { let instance = Arc::new(API{}); make_static_stream(instance) }; tokio::pin!(static_stream); while let Some(item) = static_stream.next().await { println!("Received item: {}", item); } }
方案三:用futures::stream::unfold(无额外依赖)
利用unfold将Arc<API>和原Stream绑定为状态,每次poll时推进原Stream并保留状态,确保API实例始终存活。
代码示例:
use futures::stream::{Stream, StreamExt, unfold}; use std::sync::Arc; struct API; struct Ctrl; impl API { fn get_stream(&self) -> (impl Stream<Item = i32> + '_ , Ctrl) { (futures::stream::iter(1..=5), Ctrl{}) } } fn wrap_static_stream(instance: Arc<API>) -> impl Stream<Item = i32> + 'static { let (stream, _ctrl) = instance.get_stream(); unfold((instance, stream), |(instance, mut stream)| async move { match stream.next().await { Some(item) => Some((item, (instance, stream))), None => None, } }) } #[tokio::main] async fn main() { let static_stream = { let instance = Arc::new(API{}); wrap_static_stream(instance) }; tokio::pin!(static_stream); while let Some(item) = static_stream.next().await { println!("Received item: {}", item); } }
内容的提问来源于stack exchange,提问作者JFFIGK
相关产品推荐
相关产品推荐

