AWS Lambda重复触发问题:单文件上传却多次触发DAG
问题
我创建了一个AWS Lambda,当文件上传至指定S3存储桶时触发Airflow DAG。上传文件后,Lambda会触发DAG并使其运行,但有时在DAG处理文件的过程中,即使没有新文件上传,Lambda仍会再次触发同一个DAG。
S3触发配置
Bucket arn: arn:aws:s3:::abhi-dev-backup Event types: s3:ObjectCreated:Post, s3:ObjectCreated:CompleteMultipartUpload isComplexStatement: No Prefix: com/source/trade/out/
Lambda代码
import json import boto3 import io import http.client import base64 import ast s3Client = boto3.client('s3') client = boto3.client('mwaa') #s3 = boto3.resource('s3') mwaa_env_name = 'dev-airflow' mwaa_cli_command = 'dags trigger ' def lambda_handler(event, context): dag_name = '' #Get bucket and file name bucket_name = event['Records'][0]['s3']['bucket']['name'] file_path = event['Records'][0]['s3']['object']['key'] file_name = file_path.split('/').pop() print('Bucket Name:'+bucket_name) print('File Path:'+file_path) print('File Name:'+file_name) if "SAMPLE_FILE_REQ" in file_name: dag_name = 'FILE_REQ_DAG' print('DAG NAME:'+dag_name) elif "SAMPLE_IMAGE_REQ" in file_name: dag_name = 'IMAGE_REQ_DAG' print('DAG NAME:'+dag_name) # get web token mwaa_cli_token = client.create_cli_token( Name=mwaa_env_name ) if dag_name != '': #print(mwaa_cli_token) conn = http.client.HTTPSConnection(mwaa_cli_token['WebServerHostname']) payload = mwaa_cli_command + dag_name #print(payload) headers = { 'Authorization': 'Bearer ' + mwaa_cli_token['CliToken'], 'Content-Type': 'text/plain' } conn.request("POST", "/aws_mwaa/cli", payload, headers) print('Request sent!') res = conn.getresponse() data = res.read() dict_str = data.decode("UTF-8") mydata = ast.literal_eval(dict_str) else: print('No DAG to run');
请问我哪里操作有误?为何仅上传一次文件,Lambda却会多次执行?
分析与解决方案
Lambda重复执行通常和S3事件特性、Lambda重试机制或DAG处理逻辑有关,以下是具体原因和修复方法:
1. S3事件重复触发
S3可能因网络波动、事件通知延迟重复发送同一对象创建事件;另外,若DAG处理过程中对原文件做了修改、覆盖或重新上传,会再次触发ObjectCreated类型的事件,导致Lambda重复执行。
修复:
- 调整DAG逻辑,将处理后的文件存放到其他前缀路径下,避免触发同一Lambda触发器。
- 给Lambda添加幂等性校验:比如用DynamoDB记录已处理的文件键(key),每次触发时先检查文件是否已处理,若已处理则直接返回,不触发DAG。
2. Lambda重试机制
如果Lambda执行超时、抛出未处理异常或返回错误状态码,AWS会自动重试该函数(默认最多2次)。比如你的代码中ast.literal_eval(dict_str)若解析失败会抛出异常,触发重试;或者调用MWAA的API耗时过长,导致Lambda超时中断。
修复:
- 查看CloudWatch Logs确认重复执行时的错误信息,针对性修复异常。
- 延长Lambda超时时间(比如设置为30秒以上),适配MWAA API的调用耗时。
- 在调用MWAA和解析响应的代码块添加
try-except,捕获并处理异常,避免未处理错误触发重试。
3. 事件处理逻辑缺陷
代码中直接取event['Records'][0],但S3可能批量发送多个事件,若同一文件产生多条事件记录,会导致多次触发DAG;同时没有对同一文件的重复触发做过滤。
修复:
- 遍历
event['Records']中的所有记录,对文件键进行去重处理。 - 结合幂等性校验,用文件的ETag或键作为唯一标识,确保同一文件只触发一次DAG。
示例幂等性代码片段
可以在Lambda中加入以下逻辑,用DynamoDB记录已处理文件:
# 初始化DynamoDB客户端 dynamodb = boto3.client('dynamodb') processed_table = 'processed_files' def is_file_processed(file_key): try: response = dynamodb.get_item( TableName=processed_table, Key={'file_key': {'S': file_key}} ) return 'Item' in response except Exception as e: print(f"检查处理状态出错: {e}") return False def mark_file_processed(file_key): try: dynamodb.put_item( TableName=processed_table, Item={'file_key': {'S': file_key}} ) except Exception as e: print(f"标记已处理出错: {e}") # 在lambda_handler中添加校验 if not is_file_processed(file_path): if dag_name != '': # 原有触发DAG的代码 conn = http.client.HTTPSConnection(mwaa_cli_token['WebServerHostname']) payload = mwaa_cli_command + dag_name headers = { 'Authorization': 'Bearer ' + mwaa_cli_token['CliToken'], 'Content-Type': 'text/plain' } conn.request("POST", "/aws_mwaa/cli", payload, headers) print('Request sent!') res = conn.getresponse() data = res.read() dict_str = data.decode("UTF-8") mydata = ast.literal_eval(dict_str) # 标记文件已处理 mark_file_processed(file_path) else: print(f"文件 {file_path} 已处理,跳过触发")
内容的提问来源于stack exchange,提问作者Ab_sin
相关产品推荐
相关产品推荐

