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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 04:50:54