Airflow顺序执行两次AWS Lambda时第二个任务重复触发问题
问题排查结论与修复方法
你的异常是多个代码错误+不规范用法共同导致的,核心问题如下:
- 存在拼写错误导致DAG结构解析异常:开头的DummyOperator任务把必填参数
task_id错写为task_aid,Airflow无法正确识别该任务节点,会打乱整个DAG的依赖拓扑,直接导致调度逻辑错乱,出现任务提前启动、重复触发的问题。 - 自定义boto3调用逻辑存在缺陷:
- 代码中使用的
aws_region、aws_access_key、aws_secret_access_key三个变量未做定义,会触发运行时错误,且没有配套异常捕获逻辑,容易导致任务状态判定异常 - 同步调用Lambda后未校验响应状态:就算Lambda执行超时、报错,代码没有读取响应中的
StatusCode和FunctionError字段做判断,会直接被Airflow判定为执行成功,提前触发下游任务 - 代码末尾的依赖定义行多了多余的反引号
`,会干扰DAG文件的正常解析
- 代码中使用的
- 用法不规范:没有使用Airflow官方提供的Lambda调用算子,自行实现的boto3客户端逻辑没有内置的状态校验、超时处理能力,容易出现状态判定偏差。另外Airflow 2.x版本中
provide_context参数已经废弃,不需要在default_args和算子中重复配置。
修复步骤
- 修正所有拼写、语法错误,把
task_aid改为task_id,删除依赖行末尾多余的反引号 - 替换为AWS Provider包内置的官方Lambda算子,不需要自行编写boto3逻辑,算子会自动处理同步等待、状态校验、异常抛出,保证上游任务完全执行成功后才会触发下游
- 把AWS认证信息配置到Airflow Connection中,不要硬编码在DAG代码里
修复后可运行代码
先安装必要依赖:
pip install apache-airflow-providers-amazon
DAG代码如下:
from airflow import DAG from airflow.operators.dummy import DummyOperator from datetime import timedelta from airflow.providers.amazon.aws.operators.lambda_function import LambdaInvokeFunctionOperator from airflow.utils.dates import days_ago args = { 'owner': 'airflow', 'start_date': days_ago(1), 'catchup': False } dag = DAG( 'my_dag', schedule_interval=None, default_args=args ) start = DummyOperator(task_id='Begin_execution', dag=dag) invoke_lambda1 = LambdaInvokeFunctionOperator( task_id="task1", function_name='my_lambda_function', invocation_type='RequestResponse', # 同步调用,阻塞等待Lambda执行完成再判定任务状态 payload={"date":"27/06/2022"}, execution_timeout=timedelta(hours=1), aws_conn_id='aws_default', # 提前在Airflow UI中配置AWS认证信息绑定到该连接 region_name='替换为你的AWS区域,例如ap-southeast-1', dag=dag ) invoke_lambda2 = LambdaInvokeFunctionOperator( task_id="task2", function_name='my_lambda_function', invocation_type='RequestResponse', payload={"date":"28/06/2022"}, execution_timeout=timedelta(hours=1), aws_conn_id='aws_default', region_name='替换为你的AWS区域,例如ap-southeast-1', dag=dag ) end = DummyOperator(task_id='stop_execution', dag=dag) start >> invoke_lambda1 >> invoke_lambda2 >> end
验证方法
- 部署后先执行
airflow dags list-import-errors检查DAG文件是否存在解析错误,确认无报错后再触发运行 - 运行时在Airflow UI观察任务实例状态,确认task1状态标记为
success之后,task2才会进入queued/running状态 - 核对Lambda的CloudWatch日志,确认两次调用的时间顺序、执行次数符合预期
内容的提问来源于stack exchange,提问作者mvs
相关产品推荐
相关产品推荐

