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

多文件上传后通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 11:16:11