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

Airflow顺序执行两次AWS Lambda时第二个任务重复触发问题

问题排查结论与修复方法

你的异常是多个代码错误+不规范用法共同导致的,核心问题如下:

  • 存在拼写错误导致DAG结构解析异常:开头的DummyOperator任务把必填参数task_id错写为task_aid,Airflow无法正确识别该任务节点,会打乱整个DAG的依赖拓扑,直接导致调度逻辑错乱,出现任务提前启动、重复触发的问题。
  • 自定义boto3调用逻辑存在缺陷:
    1. 代码中使用的aws_region、aws_access_key、aws_secret_access_key三个变量未做定义,会触发运行时错误,且没有配套异常捕获逻辑,容易导致任务状态判定异常
    2. 同步调用Lambda后未校验响应状态:就算Lambda执行超时、报错,代码没有读取响应中的StatusCode和FunctionError字段做判断,会直接被Airflow判定为执行成功,提前触发下游任务
    3. 代码末尾的依赖定义行多了多余的反引号`,会干扰DAG文件的正常解析
  • 用法不规范:没有使用Airflow官方提供的Lambda调用算子,自行实现的boto3客户端逻辑没有内置的状态校验、超时处理能力,容易出现状态判定偏差。另外Airflow 2.x版本中provide_context参数已经废弃,不需要在default_args和算子中重复配置。

修复步骤

  1. 修正所有拼写、语法错误,把task_aid改为task_id,删除依赖行末尾多余的反引号
  2. 替换为AWS Provider包内置的官方Lambda算子,不需要自行编写boto3逻辑,算子会自动处理同步等待、状态校验、异常抛出,保证上游任务完全执行成功后才会触发下游
  3. 把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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 14:48:19