如何在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
相关产品推荐
相关产品推荐

