如何使用Boto3将generator产出的不定长数据块写入AWS S3对象
解决方案
你当前手动实现S3多部分上传的代码存在几个核心问题:
- 不能重复使用同一个
mpu.Part()实例,每个分片需要对应唯一的分片编号,且要记录每个分片返回的ETag才能完成最终的合并操作 - S3强制要求除最后一个分片外,所有分片大小不能小于5MB,直接上传不定长的小chunk会不符合接口要求
- 原生boto3的API是同步实现,直接和异步生成器
source()混用会阻塞事件循环,异步场景推荐优先使用异步SDK
最优方案:使用SDK内置自动分片能力
不需要手动管理多部分上传逻辑,boto3/aioboto3的upload_fileobj接口会自动判断对象大小,超过阈值后自动完成分片、上传、合并的全流程,代码最简洁且不易出错:
异步场景适配(推荐)
使用异步SDK aioboto3 避免阻塞事件循环:
import aioboto3 from io import BytesIO # 替换为你的实际桶名 BUCKET = "your-bucket-name" async def write_s3(s3_key): session = aioboto3.Session() async with session.client("s3") as s3_client: byte_stream = BytesIO() # 收集异步生成器所有数据块 async for chunk in source(): byte_stream.write(chunk) # 指针重置到流开头 byte_stream.seek(0) # 内置自动处理分片上传逻辑 await s3_client.upload_fileobj(byte_stream, BUCKET, s3_key)
同步场景适配
如果不需要严格异步,也可以用原生boto3实现:
import boto3 import asyncio from io import BytesIO # 替换为你的实际桶名 BUCKET = "your-bucket-name" def write_s3(s3_key): session = boto3.Session() s3 = session.resource("s3") s3_obj = s3.Object(BUCKET, s3_key) byte_stream = BytesIO() # 迭代异步生成器收集所有数据 async def collect_chunks(): async for chunk in source(): byte_stream.write(chunk) asyncio.run(collect_chunks()) byte_stream.seek(0) s3_obj.upload_fileobj(byte_stream)
手动实现多部分上传方案
如果有自定义分片的特殊需求,可以参考以下正确实现,注意要满足S3的分片规则,且增加异常处理避免无效分片占用存储:
import boto3 import asyncio # 替换为你的实际桶名 BUCKET = "your-bucket-name" # S3强制要求非末尾分片最小为5MB MIN_PART_SIZE = 5 * 1024 * 1024 async def write_s3(s3_key): session = boto3.Session() s3 = session.resource("s3") s3_obj = s3.Object(BUCKET, s3_key) # 初始化多部分上传 mpu = s3_obj.initiate_multipart_upload() parts = [] part_number = 1 buffer = b"" try: async for chunk in source(): buffer += chunk # 攒够最小分片大小就上传 while len(buffer) >= MIN_PART_SIZE: upload_data = buffer[:MIN_PART_SIZE] buffer = buffer[MIN_PART_SIZE:] part = mpu.Part(part_number) resp = part.upload(Body=upload_data) parts.append({"PartNumber": part_number, "ETag": resp["ETag"]}) part_number += 1 # 上传最后剩余的不足5MB的分片 if buffer: part = mpu.Part(part_number) resp = part.upload(Body=buffer) parts.append({"PartNumber": part_number, "ETag": resp["ETag"]}) # 合并所有分片完成上传 mpu.complete(MultipartUpload={"Parts": parts}) except Exception as e: # 出错时终止多部分上传,清理无效分片 mpu.abort() raise e
内容的提问来源于stack exchange,提问作者ca9163d9
相关产品推荐
相关产品推荐

