Rust中如何为自定义结构体实现futures::stream::Stream trait
问题说明
需求是将reqwest获取的响应流转发给actix_web的HttpResponseBuilder,同时在流转发过程中将响应内容写入文件实现缓存,核心需要补全自定义FileCache结构体的Stream trait实现。
完整实现
首先补全结构体缺失的缓存文件句柄字段,poll_next核心逻辑为轮询上游reqwest响应流,拿到数据块后先写入缓存文件,再将数据透传给下游actix,流结束时刷盘保证缓存完整。
use std::pin::Pin; use std::task::{Context, Poll}; use bytes::Bytes; use futures::Stream; use tokio::fs::File; use tokio::io::AsyncWriteExt; struct FileCache { // 上游reqwest响应流,加Send约束满足actix线程调度要求 stream: Pin<Box<dyn Stream<Item = reqwest::Result<Bytes>> + Send>>, // 缓存文件异步写句柄 cache_file: File, } impl FileCache { fn new( stream: Box<dyn Stream<Item = reqwest::Result<Bytes>> + Send>, cache_file: File ) -> Self { Self { stream: Pin::new(stream), cache_file, } } } impl Stream for FileCache { type Item = reqwest::Result<Bytes>; fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> { let this = self.as_mut().get_mut(); match this.stream.as_mut().poll_next(cx) { // 上游还没准备好数据,直接透传Pending状态 Poll::Pending => Poll::Pending, // 上游流传输结束,刷盘缓存文件后返回结束标识 Poll::Ready(None) => { let _ = std::task::ready!(Pin::new(&mut this.cache_file).poll_flush(cx)); Poll::Ready(None) } // 上游返回错误,直接透传 Poll::Ready(Some(Err(e))) => Poll::Ready(Some(Err(e))), // 拿到正常数据块 Poll::Ready(Some(Ok(bytes))) => { // 异步写入缓存文件,不阻塞事件循环 match std::task::ready!(Pin::new(&mut this.cache_file).poll_write(cx, &bytes)) { Ok(write_len) => { if write_len != bytes.len() { // *缓存是附加功能,写异常建议打日志即可,不要中断主响应流程* eprintln!("cache write length mismatch, expect {} got {}", bytes.len(), write_len); } } Err(e) => { eprintln!("cache write failed: {}", e); } } // 原始数据透传给下游actix响应 Poll::Ready(Some(Ok(bytes))) } } } }
使用方式
在actix handler中直接包装reqwest的字节流传入即可,全程逐块处理不会将全量响应加载到内存:
use actix_web::{HttpResponse, get}; use reqwest::Client; #[get("/proxy")] async fn proxy() -> HttpResponse { // 发起请求拿到reqwest响应 let upstream_resp = Client::new() .get("https://target-domain.com/resource") .send() .await .unwrap(); // 提前创建缓存文件,实际场景注意处理路径创建、文件名去重等逻辑 let cache_file = File::create("./local-cache/resource.bin").await.unwrap(); // 包装成缓存流 let cache_stream = FileCache::new(Box::new(upstream_resp.bytes_stream()), cache_file); // 传给actix构建流式响应 HttpResponse::Ok() .content_type(upstream_resp.headers().get("content-type").unwrap().to_str().unwrap()) .streaming(cache_stream) }
注意事项
- 不要在
poll_next中使用阻塞IO(比如std::fs::File的同步写方法),会阻塞tokio事件循环降低服务吞吐量,必须使用异步IO或者把阻塞写逻辑放到spawn_blocking中执行 - 给上游流加
Send约束是必须的,否则actix跨线程调度流时会报编译错误 - 生产环境建议给缓存逻辑加超时、异常降级逻辑,不要因为缓存组件故障影响正常接口响应
- 可以根据需求扩展逻辑,比如根据响应头判断是否需要缓存、设置缓存过期时间等
内容的提问来源于stack exchange,提问作者Dr. Light
相关产品推荐
相关产品推荐

