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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 01:57:22