利用boto3实现S3对象存储的高效分片上传与断点续传
1. 为什么需要分片上传与断点续传
当你需要上传一个10GB的视频文件到S3对象存储时,直接使用普通上传会遇到几个头疼的问题。首先是上传时间长,网络稍微波动就可能前功尽弃;其次是内存占用高,Python进程可能直接崩溃;最重要的是,一旦中断就得从头再来。
我去年就遇到过这种情况:客户现场的网络环境不稳定,上传一个8GB的数据库备份文件时,每次到90%左右就断开连接。重复尝试了三次都没成功,最后只能带着硬盘去客户现场物理拷贝。这就是为什么我们需要掌握分片上传和断点续传这两个核心技术。
分片上传的原理很简单:把大文件切成小块(比如每块100MB),然后并行上传这些小块。即使某个分片上传失败,也只需要重传这个分片,而不是整个文件。实际测试中,对一个5GB的文件使用分片上传,速度比普通上传快3倍以上,而且内存占用稳定在200MB以内。
2. 配置boto3环境与基础操作
2.1 安装与基础配置
首先确保你已经安装了boto3库。如果还没安装,用pip一键搞定:
pip install boto3
配置认证信息时,我强烈建议不要直接把AK/SK硬编码在代码里。更好的做法是使用AWS CLI配置凭证文件:
aws configure
然后在代码中这样初始化client:
import boto3
s3 = boto3.client('s3',
endpoint_url='https://your-s3-endpoint.com',
region_name='us-east-1')
如果是非AWS的S3兼容存储(比如阿里云OSS、腾讯云COS),需要额外指定签名版本:
from botocore.client import Config
s3 = boto3.client('s3',
endpoint_url='https://your-endpoint.com',
config=Config(signature_version='s3v4'),
region_name='your-region')
2.2 基础操作验证
上传前先做个简单的连通性测试:
# 列出所有bucket测试连接
response = s3.list_buckets()
print([bucket['Name'] for bucket in response['Buckets']])
# 上传测试文件
with open('test.txt', 'w') as f:
f.write('hello world')
s3.upload_file('test.txt', 'your-bucket', 'test.txt')
如果这些基础操作都正常,说明环境配置没问题。遇到过证书问题的同学可以加上verify=False参数,但生产环境建议正确配置证书。
3. 实现高效分片上传
3.1 分片上传全流程
完整的分片上传包含三个关键步骤:
- 初始化上传:获取唯一的Upload ID
- 上传分片:并行上传各个分片
- 完成上传:合并所有分片
来看具体代码实现:
def multipart_upload(file_path, bucket, object_key, chunk_size=100*1024*1024):
# 1. 初始化
upload_id = s3.create_multipart_upload(
Bucket=bucket,
Key=object_key
)['UploadId']
parts = []
try:
# 2. 分片上传
with open(file_path, 'rb') as f:
part_number = 1
while True:
chunk = f.read(chunk_size)
if not chunk:
break
response = s3.upload_part(
Bucket=bucket,
Key=object_key,
PartNumber=part_number,
UploadId=upload_id,
Body=chunk
)
parts.append({
'PartNumber': part_number,
'ETag': response['ETag']
})
part_number += 1
# 3. 完成上传
s3.complete_multipart_upload(
Bucket=bucket,
Key=object_key,
UploadId=upload_id,
MultipartUpload={'Parts': parts}
)
print("上传成功")
except Exception as e:
print(f"上传失败: {e}")
s3.abort_multipart_upload(
Bucket=bucket,
Key=object_key,
UploadId=upload_id
)
raise
3.2 性能优化技巧
通过以下几个技巧可以显著提升上传速度:
- 并行上传:使用线程池同时上传多个分片
- 动态分片大小:根据网络质量调整分片大小
- 内存优化:使用生成器避免大文件一次性加载
改进后的并行上传代码:
from concurrent.futures import ThreadPoolExecutor
def upload_part(args):
part_num, chunk, upload_id = args
response = s3.upload_part(
Bucket=bucket,
Key=object_key,
PartNumber=part_num,
UploadId=upload_id,
Body=chunk
)
return {'PartNumber': part_num, 'ETag': response['ETag']}
def parallel_upload(file_path, bucket, object_key, workers=4):
upload_id = s3.create_multipart_upload(Bucket=bucket, Key=object_key)['UploadId']
parts = []
try:
with ThreadPoolExecutor(max_workers=workers) as executor:
futures = []
part_num = 1
with open(file_path, 'rb') as f:
while True:
chunk = f.read(100*1024*1024) # 100MB
if not chunk:
break
futures.append(executor.submit(
upload_part,
(part_num, chunk, upload_id)
))
part_num += 1
for future in futures:
parts.append(future.result())
s3.complete_multipart_upload(
Bucket=bucket,
Key=object_key,
UploadId=upload_id,
MultipartUpload={'Parts': sorted(parts, key=lambda x: x['PartNumber'])}
)
except Exception:
s3.abort_multipart_upload(Bucket=bucket, Key=object_key, UploadId=upload_id)
raise
实测显示,4线程并行上传5GB文件,耗时从单线程的8分钟降至2分15秒。
4. 实现断点续传机制
4.1 断点续传原理
断点续传需要解决三个核心问题:
- 记录上传进度:保存已上传的分片信息
- 恢复上传:从断点处继续上传
- 处理异常:网络中断后的清理工作
4.2 完整实现方案
import os
import pickle
from datetime import datetime
def resume_upload(file_path, bucket, object_key, state_file='upload.state'):
# 尝试加载进度
if os.path.exists(state_file):
with open(state_file, 'rb') as f:
state = pickle.load(f)
upload_id = state['upload_id']
parts = state['parts']
else:
upload_id = s3.create_multipart_upload(Bucket=bucket, Key=object_key)['UploadId']
parts = []
try:
file_size = os.path.getsize(file_path)
chunk_size = 100 * 1024 * 1024 # 100MB
total_parts = (file_size + chunk_size - 1) // chunk_size
with open(file_path, 'rb') as f:
for part_num in range(1, total_parts + 1):
# 跳过已上传的分片
if any(p['PartNumber'] == part_num for p in parts):
f.seek(chunk_size * (part_num - 1))
continue
f.seek(chunk_size * (part_num - 1))
chunk = f.read(chunk_size)
response = s3.upload_part(
Bucket=bucket,
Key=object_key,
PartNumber=part_num,
UploadId=upload_id,
Body=chunk
)
parts.append({
'PartNumber': part_num,
'ETag': response['ETag']
})
# 保存进度
with open(state_file, 'wb') as f_state:
pickle.dump({
'upload_id': upload_id,
'parts': parts,
'last_updated': datetime.now().isoformat()
}, f_state)
# 完成上传
s3.complete_multipart_upload(
Bucket=bucket,
Key=object_key,
UploadId=upload_id,
MultipartUpload={'Parts': sorted(parts, key=lambda x: x['PartNumber'])}
)
os.remove(state_file) # 清理状态文件
print("上传完成")
except Exception as e:
print(f"上传中断: {e}")
print("下次运行将自动从断点继续")
raise
这个实现会在本地保存一个状态文件,记录已上传的分片信息。即使程序崩溃或网络中断,再次运行时会自动从上次中断的位置继续上传。
5. 实战中的常见问题与解决方案
5.1 分片大小选择
分片大小直接影响上传性能:
- 太小:增加管理开销,降低速度
- 太大:重传成本高,内存压力大
建议值:
- 高速稳定网络:100-200MB
- 普通网络:50-100MB
- 不稳定网络:10-50MB
可以通过以下代码动态调整分片大小:
def get_optimal_chunk_size(file_size, network_speed_mbps=10):
"""根据文件大小和网络状况计算最佳分片大小"""
min_chunk = 5 * 1024 * 1024 # 5MB最小
max_chunk = 500 * 1024 * 1024 # 500MB最大
# 每MB需要的上传时间(秒)
time_per_mb = 8 / network_speed_mbps
# 目标:每个分片上传时间在30-60秒之间
ideal_chunk = (30 / time_per_mb) * 1024 * 1024
return min(max_chunk, max(min_chunk, ideal_chunk))
5.2 错误处理与重试机制
网络不稳定时,需要实现智能重试:
from botocore.exceptions import ClientError
import time
def upload_part_with_retry(bucket, key, part_num, chunk, upload_id, max_retries=3):
for attempt in range(max_retries):
try:
response = s3.upload_part(
Bucket=bucket,
Key=key,
PartNumber=part_num,
UploadId=upload_id,
Body=chunk
)
return response
except ClientError as e:
if attempt == max_retries - 1:
raise
wait_time = (2 ** attempt) * 0.5 # 指数退避
time.sleep(wait_time)
5.3 跨区域上传优化
当客户端与S3存储不在同一区域时,上传速度会明显下降。解决方法:
- 使用传输加速端点:
s3 = boto3.client('s3',
endpoint_url='https://s3-accelerate.amazonaws.com',
config=Config(signature_version='s3v4'))
-
通过CDN边缘节点上传
-
使用AWS Direct Connect专线
6. 高级应用场景
6.1 大文件下载的断点续传
下载同样需要断点续传机制:
def resume_download(bucket, key, file_path, chunk_size=100*1024*1024):
downloaded = 0
if os.path.exists(file_path):
downloaded = os.path.getsize(file_path)
headers = {'Range': f'bytes={downloaded}-'}
response = s3.get_object(Bucket=bucket, Key=key, **headers)
with open(file_path, 'ab') as f: # 追加模式
for chunk in response['Body'].iter_chunks(chunk_size):
f.write(chunk)
6.2 与Lambda函数结合
在无服务器架构中处理大文件上传:
import boto3
from urllib.parse import parse_qs
def lambda_handler(event, context):
s3 = boto3.client('s3')
params = parse_qs(event['body'])
# 处理分片上传
if 'uploadId' in params:
# 继续上传逻辑
pass
else:
# 初始化上传
response = s3.create_multipart_upload(
Bucket=params['bucket'][0],
Key=params['key'][0]
)
return {
'statusCode': 200,
'body': response['UploadId']
}
6.3 客户端加密上传
对于敏感数据,可以在客户端加密后再上传:
from cryptography.fernet import Fernet
def encrypted_upload(file_path, bucket, key):
# 生成加密密钥
key = Fernet.generate_key()
cipher = Fernet(key)
with open(file_path, 'rb') as f:
encrypted_data = cipher.encrypt(f.read())
s3.put_object(
Bucket=bucket,
Key=key,
Body=encrypted_data,
Metadata={
'encryption-key': key.decode('utf-8') # 实际项目应该用KMS管理密钥
}
)
7. 监控与性能调优
7.1 上传进度监控
添加进度条显示上传进度:
from tqdm import tqdm
def upload_with_progress(file_path, bucket, key):
file_size = os.path.getsize(file_path)
with tqdm(total=file_size, unit='B', unit_scale=True) as pbar:
def callback(bytes_transferred):
pbar.update(bytes_transferred)
s3.upload_file(
Filename=file_path,
Bucket=bucket,
Key=key,
Callback=callback
)
7.2 S3性能指标分析
通过CloudWatch监控上传指标:
cloudwatch = boto3.client('cloudwatch')
def get_upload_metrics(bucket):
response = cloudwatch.get_metric_statistics(
Namespace='AWS/S3',
MetricName='BytesUploaded',
Dimensions=[{'Name': 'BucketName', 'Value': bucket}],
StartTime=datetime.utcnow() - timedelta(days=1),
EndTime=datetime.utcnow(),
Period=3600,
Statistics=['Sum']
)
return response['Datapoints']
7.3 成本优化建议
- 对于不常访问的数据,使用S3 Infrequent Access存储类
- 设置生命周期策略自动转移旧数据
- 对已完成上传的分片及时调用complete操作,避免产生不必要的存储费用
# 设置生命周期策略
lifecycle = {
'Rules': [
{
'ID': 'Move to IA after 30 days',
'Status': 'Enabled',
'Prefix': '',
'Transitions': [
{
'Days': 30,
'StorageClass': 'STANDARD_IA'
}
]
}
]
}
s3.put_bucket_lifecycle_configuration(
Bucket='your-bucket',
LifecycleConfiguration=lifecycle
)
8. 安全最佳实践
8.1 权限控制
遵循最小权限原则,使用IAM策略限制上传权限:
{
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Action": [
"s3:PutObject",
"s3:AbortMultipartUpload",
"s3:ListMultipartUploadParts"
],
"Resource": "arn:aws:s3:::your-bucket/*"
}
]
}
8.2 数据完整性校验
上传完成后验证ETag确保数据完整:
def verify_upload(bucket, key, local_file):
# 获取本地文件MD5
with open(local_file, 'rb') as f:
local_md5 = hashlib.md5(f.read()).hexdigest()
# 获取S3对象ETag
response = s3.head_object(Bucket=bucket, Key=key)
etag = response['ETag'].strip('"')
if local_md5 == etag:
print("文件校验通过")
else:
print("文件校验失败")
raise ValueError("文件内容不匹配")
8.3 临时凭证使用
使用STS获取临时凭证进行上传:
sts = boto3.client('sts')
response = sts.assume_role(
RoleArn='arn:aws:iam::123456789012:role/UploadRole',
RoleSessionName='upload-session'
)
credentials = response['Credentials']
s3 = boto3.client(
's3',
aws_access_key_id=credentials['AccessKeyId'],
aws_secret_access_key=credentials['SecretAccessKey'],
aws_session_token=credentials['SessionToken']
)
更多推荐


所有评论(0)