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

如何使用Boto3将generator产出的不定长数据块写入AWS S3对象

解决方案

你当前手动实现S3多部分上传的代码存在几个核心问题:

  1. 不能重复使用同一个mpu.Part()实例,每个分片需要对应唯一的分片编号,且要记录每个分片返回的ETag才能完成最终的合并操作
  2. S3强制要求除最后一个分片外,所有分片大小不能小于5MB,直接上传不定长的小chunk会不符合接口要求
  3. 原生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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 11:15:07