如何通过Lambda传入配置触发Airflow动态DAG
实现方案
整个流程分两步走:先保证动态DAG脚本能基于Lambda传入的配置正确生成可被Airflow识别的DAG,再通过Airflow原生开放接口实现Lambda侧的指定DAG触发。
一、动态DAG脚本的加载与更新逻辑
Airflow的DAG不是你代码里运行生成个对象就能直接调度的,必须被Scheduler扫描解析、同步到元数据库后才能被触发,所以不能直接让Lambda把配置传给DAG脚本做即时运行,得走持久化配置+Airflow自动/手动解析的标准流程,别搞Lambda直接跑DAG脚本的骚操作,没用:
- 先选一个Airflow和Lambda都能读写的持久化存储放配置,可选Airflow自带Variable、S3存储JSON配置文件、业务数据库配置表都行,Lambda每次调整工作流配置,直接把最新的全量DAG配置列表写入这个存储位置即可,配置结构里要给每个工作流分配全局唯一的
dag_id。 - 把你写的动态DAG生成脚本放到Airflow配置的
dags_folder目录下,脚本核心逻辑是启动时读取存储里的配置,循环遍历每个配置项生成独立的DAG对象,最后必须把生成的DAG对象注册到全局命名空间,不然Airflow扫描不到。参考代码如下:
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime import json # 替换成你实际的配置读取逻辑,比如读S3、读数据库、读Airflow Variable def load_dynamic_configs(): from airflow.models import Variable return json.loads(Variable.get("dynamic_workflow_configs", default_var="[]")) all_configs = load_dynamic_configs() for config in all_configs: current_dag_id = config["dag_id"] with DAG( dag_id=current_dag_id, start_date=datetime(2024, 1, 1), schedule=None, catchup=False, default_args={"owner": "dynamic_workflow"} ) as dag: # 基于config里的自定义参数定义具体任务逻辑 PythonOperator( task_id="execute_workflow", python_callable=lambda **ctx: print(f"执行工作流{current_dag_id},入参:{config}") ) # 关键操作:将生成的DAG注册到全局作用域,Airflow才能识别 globals()[current_dag_id] = dag
- 配置更新后的DAG刷新:Airflow默认会每隔
dag_dir_list_interval(默认300秒)扫描一次DAG目录,非实时场景等自动扫描即可;如果要配置更新后立刻生效,Lambda写完配置后可以直接调用Airflow REST API触发DAG重新解析,不需要重启Airflow服务。
二、Lambda端触发指定生成DAG的操作
等Airflow解析出对应dag_id的DAG后,直接走Airflow官方REST API触发即可,不要尝试在Lambda侧直接加载DAG文件运行,会出现元数据不一致的问题:
- 提前在Airflow里创建专用的服务账号,给这个账号分配目标DAG的触发权限,不要用管理员账号做接口调用,配置好Token或者Basic Auth认证信息存到Lambda的环境变量里。
- Lambda确认配置同步完成、DAG解析成功后,调用触发DAG运行的接口:接口路径为
POST /api/v1/dags/{你要触发的目标DAG的dag_id}/dagRuns,请求体可以传入自定义运行参数,示例:
{ "logical_date": "2024-05-20T00:00:00Z", "conf": { "custom_biz_param": "test_value" }, "note": "triggered from lambda" }
- 实操避坑点:
- 配置里的
dag_id必须全局唯一,不能和已有的静态DAG重名,否则会覆盖原有DAG导致调度异常 - Lambda调用触发接口时建议加2-3次重试,每次间隔10秒,避免刚写完配置Airflow还没完成DAG解析就调用,触发“DAG不存在”的报错
- 不要在Lambda里直接操作Airflow元数据库写DAG运行记录,所有操作走官方REST API,避免版本升级后元数据结构不兼容
- 别信网上那种在Lambda里手动加载DagBag生成DAG的野路子,生产环境很容易把元数据库搞乱,排查问题成本极高
- 配置里的
内容的提问来源于stack exchange,提问作者anukriti kulshrestha
相关产品推荐
相关产品推荐

