多文件上传后通过SNS Topic触发Lambda的实现方案
解决方案:S3多文件上传完成后触发Lambda并通知SNS
核心思路
通过DynamoDB存储批次上传状态,结合S3事件触发的Lambda实现文件计数,当判断批次内所有文件上传完成后,执行业务逻辑并触发SNS Topic。针对文件数量不固定的场景,提供两种可行的批次判断方案:
方案1:基于Manifest文件的固定数量判断
适合客户可以提前知晓上传文件总数的场景,流程如下:
1. 准备工作
- 创建DynamoDB表:
- 分区键:
BatchID(字符串类型,用于标识同一批次的文件,比如客户自定义的批次ID) - 字段:
TotalExpected(数字,该批次总文件数)、UploadedCount(数字,已上传文件数)、Processed(布尔,标记批次是否已处理完成)
- 分区键:
2. S3事件配置
给目标S3桶添加事件通知:
- 事件类型:
对象创建(所有) - 触发目标:指定你的Lambda函数
- 可选过滤规则:如果用文件夹区分批次,可设置前缀(如
batch-123/)缩小触发范围
3. Lambda逻辑实现
import boto3 import json dynamodb = boto3.resource('dynamodb') sns = boto3.client('sns') table = dynamodb.Table('BatchUploadStatus') SNS_TOPIC_ARN = 'arn:aws:sns:region:account-id:your-topic' def lambda_handler(event, context): # 解析S3事件中的对象信息 s3_obj = event['Records'][0]['s3'] obj_key = s3_obj['object']['key'] # 处理manifest文件(客户先上传该文件,包含批次ID和总文件数) if obj_key.endswith('manifest.json'): s3_client = boto3.client('s3') manifest_content = s3_client.get_object(Bucket='your-bucket', Key=obj_key)['Body'].read().decode() manifest = json.loads(manifest_content) # 初始化DynamoDB中的批次状态 table.put_item( Item={ 'BatchID': manifest['batch_id'], 'TotalExpected': manifest['total_files'], 'UploadedCount': 0, 'Processed': False } ) return # 处理普通上传文件 batch_id = obj_key.split('/')[0] # 假设用文件夹前缀作为批次ID # 幂等更新:避免S3重复触发导致计数错误 response = table.update_item( Key={'BatchID': batch_id}, UpdateExpression='SET UploadedCount = UploadedCount + :inc', ConditionExpression='Processed = :false', ExpressionAttributeValues={':inc': 1, ':false': False}, ReturnValues='ALL_NEW' ) updated_item = response['Attributes'] # 判断是否所有文件上传完成 if updated_item['UploadedCount'] == updated_item['TotalExpected']: # 标记批次为已处理 table.update_item( Key={'BatchID': batch_id}, UpdateExpression='SET Processed = :true', ExpressionAttributeValues={':true': True} ) # 执行你的业务逻辑(根据文件总数返回不同结果) # ... # 通知SNS Topic sns.publish( TopicArn=SNS_TOPIC_ARN, Message=json.dumps({ 'batch_id': batch_id, 'total_files': updated_item['TotalExpected'], 'status': 'completed' }), Subject='Batch Upload Completed' )
方案2:基于超时的动态数量判断
适合无法提前知晓文件总数的场景,通过设置超时时间(比如5分钟)判断批次是否完成:
1. 准备工作
- 同上创建DynamoDB表,新增字段
LastUploadTime(时间戳,记录该批次最后一次文件上传时间) - 创建CloudWatch Events规则:定时触发Lambda(比如每分钟一次),用于检查超时未更新的批次
2. Lambda逻辑调整
- 文件上传触发的Lambda:更新对应批次的
UploadedCount和LastUploadTime - 定时Lambda:查询DynamoDB中
Processed=False且LastUploadTime早于当前时间-超时时间的批次,判定为上传完成,执行业务逻辑并触发SNS
权限策略配置
1. Lambda执行角色权限(IAM Policy)
{ "Version": "2012-10-17", "Statement": [ { "Effect": "Allow", "Action": [ "s3:GetObject", "s3:ListBucket" ], "Resource": [ "arn:aws:s3:::your-bucket", "arn:aws:s3:::your-bucket/*" ] }, { "Effect": "Allow", "Action": [ "dynamodb:GetItem", "dynamodb:PutItem", "dynamodb:UpdateItem", "dynamodb:Scan" ], "Resource": "arn:aws:dynamodb:region:account-id:table/BatchUploadStatus" }, { "Effect": "Allow", "Action": "sns:Publish", "Resource": "arn:aws:sns:region:account-id:your-topic" }, { "Effect": "Allow", "Action": [ "logs:CreateLogGroup", "logs:CreateLogStream", "logs:PutLogEvents" ], "Resource": "arn:aws:logs:*:*:*" } ] }
2. S3调用Lambda的权限(Lambda资源策略)
确保S3有权触发Lambda,在Lambda的资源策略中添加:
{ "Version": "2012-10-17", "Id": "s3-event-trigger", "Statement": [ { "Sid": "AllowS3ToInvokeLambda", "Effect": "Allow", "Principal": { "Service": "s3.amazonaws.com" }, "Action": "lambda:InvokeFunction", "Resource": "arn:aws:lambda:region:account-id:function:your-lambda-function", "Condition": { "ArnLike": { "AWS:SourceArn": "arn:aws:s3:::your-bucket" } } } ] }
关键注意事项
- 幂等性:S3事件可能重复触发,Lambda逻辑要通过DynamoDB条件更新、记录已处理文件ID等方式避免重复计数
- 数据清理:定期清理DynamoDB中已处理的批次数据,避免存储冗余
- 超时设置:方案2中的超时时间需根据实际上传速度调整,避免误判
内容的提问来源于stack exchange,提问作者clattenburg cake
相关产品推荐
相关产品推荐

