如何使用Boto3确认S3大文件写入完成后再启动跨桶复制
确保S3大文件完全写入后再触发跨桶复制的可靠方案
Great question—handling large file writes (5GB+) in S3 and making sure you only copy them once they’re fully uploaded is a common pain point in distributed workflows. Let’s break down the most reliable approaches, with Boto3 code examples you can use directly:
方案1:利用S3 ETag进行完整性验证
S3的ETag(实体标签)是判断文件是否完整上传的关键,尤其是对于分块上传的大文件(5GB+必须用分块上传,因为S3的单PUT请求最大支持5GB)。
原理:
- 对于单PUT上传(≤5GB):ETag就是文件的MD5哈希值
- 对于分块上传(>5GB):ETag是每个分块MD5哈希的拼接,再加上分块数量的后缀(格式类似
abc123-def,其中def是分块数) - 进程A完成上传后,可以预先计算好正确的ETag;进程B先获取S3上对象的ETag,对比一致后再执行复制。
Boto3代码示例:
import boto3 import hashlib from math import ceil s3_client = boto3.client('s3') def calculate_multipart_etag(file_path, chunk_size=8*1024*1024): """计算分块上传文件的ETag(默认块大小8MB,和S3默认分块大小一致)""" md5s = [] with open(file_path, 'rb') as f: while chunk := f.read(chunk_size): md5s.append(hashlib.md5(chunk)) if len(md5s) == 1: return f"{md5s[0].hexdigest()}" else: combined_md5 = hashlib.md5(b''.join(m.digest() for m in md5s)) return f"{combined_md5.hexdigest()}-{len(md5s)}" def copy_file_after_verification(source_bucket, source_key, dest_bucket, dest_key, expected_etag): # 获取S3上对象的ETag response = s3_client.head_object(Bucket=source_bucket, Key=source_key) s3_etag = response['ETag'].strip('"') # S3返回的ETag带双引号,需要去掉 if s3_etag == expected_etag: print("文件完整性验证通过,开始复制...") s3_client.copy_object( Bucket=dest_bucket, Key=dest_key, CopySource={'Bucket': source_bucket, 'Key': source_key} ) print("复制完成") else: print("文件未完全上传或损坏,跳过复制") # 使用示例:进程A上传完后把计算好的expected_etag传给进程B expected_etag = calculate_multipart_etag('/path/to/large/file') copy_file_after_verification('bucket-a', 'large-file.dat', 'bucket-b', 'copied-large-file.dat', expected_etag)
方案2:使用S3事件通知触发复制(自动化首选)
如果你的流程是异步的,最省心的方式是让S3在对象完全创建完成后自动触发进程B(或者Lambda函数)执行复制。
原理:
- 给Bucket A配置S3事件通知,选择
s3:ObjectCreated:*事件(包含分块上传完成的s3:ObjectCreated:CompleteMultipartUpload) - 可以把事件发送到Lambda、SQS或者SNS,进程B监听这些服务的消息,收到后直接执行复制操作。
Boto3配置事件通知示例(以Lambda为目标):
def configure_s3_event_notification(bucket_name, lambda_arn): s3_client.put_bucket_notification_configuration( Bucket=bucket_name, NotificationConfiguration={ 'LambdaFunctionConfigurations': [ { 'LambdaFunctionArn': lambda_arn, 'Events': ['s3:ObjectCreated:CompleteMultipartUpload', 's3:ObjectCreated:Put'], 'Filter': { 'Key': { 'FilterRules': [ { 'Name': 'size', 'Value': '5242880000' # 5GB,只触发大文件的事件 } ] } } } ] } )
然后在Lambda函数里写复制逻辑:
import boto3 s3_client = boto3.client('s3') def lambda_handler(event, context): # 从事件中获取源对象信息 record = event['Records'][0] source_bucket = record['s3']['bucket']['name'] source_key = record['s3']['object']['key'] # 复制到Bucket B dest_bucket = 'bucket-b' dest_key = f"copied/{source_key}" s3_client.copy_object( Bucket=dest_bucket, Key=dest_key, CopySource={'Bucket': source_bucket, 'Key': source_key} ) return f"Successfully copied {source_key} to {dest_bucket}"
方案3:使用“完成标记”文件(简单场景首选)
如果你的架构比较简单,不想搞复杂的事件或ETag计算,可以让进程A在完全上传大文件后,上传一个小的“标记文件”(比如large-file.dat.completed)。进程B轮询这个标记文件,只有当标记存在时才复制大文件。
Boto3代码示例:
import time def wait_for_completion_mark(source_bucket, source_key, poll_interval=10): """轮询等待完成标记文件""" completion_key = f"{source_key}.completed" while True: try: s3_client.head_object(Bucket=source_bucket, Key=completion_key) print("检测到完成标记,开始复制") return True except s3_client.exceptions.ClientError as e: if e.response['Error']['Code'] == '404': print("完成标记未出现,等待中...") time.sleep(poll_interval) else: raise def copy_file_with_mark(source_bucket, source_key, dest_bucket, dest_key): if wait_for_completion_mark(source_bucket, source_key): s3_client.copy_object( Bucket=dest_bucket, Key=dest_key, CopySource={'Bucket': source_bucket, 'Key': source_key} ) print("复制完成") # 使用示例 copy_file_with_mark('bucket-a', 'large-file.dat', 'bucket-b', 'copied-large-file.dat')
注意事项
- 不管用哪种方案,进程A必须确保大文件的上传是原子完成的:如果用分块上传,一定要调用
complete_multipart_upload,否则S3不会生成完整对象;如果用单PUT,要确保请求成功返回。 - 跨桶复制需要确保进程B(或Lambda)有足够的IAM权限:
s3:GetObject(对Bucket A)和s3:PutObject(对Bucket B)。
内容的提问来源于stack exchange,提问作者sandy
相关产品推荐
相关产品推荐

