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

如何在Axum中通过Multipart流式上传文件至S3(rust-s3)

实现Axum Multipart Field到S3 AsyncRead的零拷贝流式上传

要实现Axum的multipart::Field和rust-s3的put_object_stream协同工作,核心是把Field转换成tokio::io::AsyncRead类型,这样就能直接流式上传到S3,无需将整个文件缓存到内存或磁盘。

步骤1:定义包装结构体

先创建一个结构体来包装axum::extract::multipart::Field,后续为它实现AsyncRead trait:

use axum::extract::multipart::Field;
use tokio::io::{AsyncRead, ReadBuf};
use std::pin::Pin;
use std::task::{Context, Poll};

struct FieldAsyncRead {
    field: Field,
    current_chunk: Option<bytes::Bytes>,
    chunk_pos: usize,
}

impl FieldAsyncRead {
    fn new(field: Field) -> Self {
        Self {
            field,
            current_chunk: None,
            chunk_pos: 0,
        }
    }
}

步骤2:实现AsyncRead trait

为FieldAsyncRead实现AsyncRead,核心逻辑是:当前chunk的数据读完后,异步获取下一个chunk,再继续填充目标缓冲区:

impl AsyncRead for FieldAsyncRead {
    fn poll_read(
        mut self: Pin<&mut Self>,
        cx: &mut Context<'_>,
        buf: &mut ReadBuf<'_>,
    ) -> Poll<std::io::Result<()>> {
        loop {
            // 优先处理当前未读完的chunk
            if let Some(chunk) = &mut self.current_chunk {
                let remaining = &chunk[self.chunk_pos..];
                let copy_len = remaining.len().min(buf.remaining());
                
                buf.put_slice(&remaining[..copy_len]);
                self.chunk_pos += copy_len;
                
                // 当前chunk读完则清空,准备读取下一个
                if self.chunk_pos == chunk.len() {
                    self.current_chunk = None;
                    self.chunk_pos = 0;
                }
                
                // 有数据拷贝完成就返回Ready
                if copy_len > 0 {
                    return Poll::Ready(Ok(()));
                }
            }
            
            // 无可用chunk时,尝试获取下一个
            match Pin::new(&mut self.field).poll_chunk(cx) {
                Poll::Ready(Ok(Some(chunk))) => {
                    self.current_chunk = Some(chunk);
                    self.chunk_pos = 0;
                }
                Poll::Ready(Ok(None)) => {
                    // 所有chunk读取完毕,返回EOF
                    return Poll::Ready(Ok(()));
                }
                Poll::Ready(Err(e)) => {
                    // 将Axum错误转换为std::io::Error
                    let io_err = std::io::Error::new(
                        std::io::ErrorKind::Other,
                        format!("读取Multipart块失败: {}", e)
                    );
                    return Poll::Ready(Err(io_err));
                }
                Poll::Pending => {
                    // 未获取到新chunk,返回Pending等待调度
                    return Poll::Pending;
                }
            }
        }
    }
}

步骤3:在Axum端点中使用

在你的多部分表单处理逻辑里,把Field转换成FieldAsyncRead,直接传给s3::Bucket::put_object_stream即可:

use axum::{extract::Multipart, routing::post, Router};
use s3::Bucket;
use std::sync::Arc;

async fn upload_handler(mut multipart: Multipart, bucket: Arc<Bucket>) -> Result<String, String> {
    while let Some(field) = multipart.next_field().await.map_err(|e| e.to_string())? {
        let filename = field.file_name().ok_or("缺少文件名")?.to_string();
        
        // 将Field包装为AsyncRead类型
        let async_read = FieldAsyncRead::new(field);
        
        // 流式上传到S3
        bucket
            .put_object_stream(async_read, filename.as_str())
            .await
            .map_err(|e| e.to_string())?;
    }
    
    Ok("上传完成".to_string())
}

// 路由示例
#[tokio::main]
async fn main() {
    // 初始化S3 Bucket(根据实际配置调整)
    let bucket = Arc::new(s3::Bucket::new(
        "你的存储桶名称",
        s3::Region::UsEast1,
        s3::creds::Credentials::default().unwrap(),
    ).unwrap());
    
    let app = Router::new()
        .route("/upload", post(upload_handler))
        .with_state(bucket);
    
    axum::Server::bind(&"0.0.0.0:3000".parse().unwrap())
        .serve(app.into_make_svc())
        .await
        .unwrap();
}

关键说明

  • 全程流式处理:每读取一个Multipart块就立即传给S3,不会缓存整个文件到内存或磁盘。
  • 错误类型转换:将Axum的Multipart错误转换为std::io::Error,满足AsyncRead的错误类型要求。
  • 依赖配置:确保rust-s3启用tokio特性,在Cargo.toml中配置:s3 = { version = "0.34", features = ["tokio"] }

内容的提问来源于stack exchange,提问作者kiv_apple

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 09:57:17