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

Airflow:如何在非调度器管控代码中发布并触发动态DAG

嘿,我完全懂你不想把DAG硬塞进$AIRFLOW_HOME/dags目录的困扰——生成py文件的方式确实有点绕。既然你已经在dag_pickle表看到了条目,说明你已经摸到门道了,接下来咱们把剩下的步骤理清楚:

优先推荐:用Airflow REST API实现无文件发布+触发

这应该是你要的「更简便方法」,完全不需要生成py文件,直接在外部代码里完成DAG发布和触发,全程和调度器的dags目录无关:

  • 第一步:发布DAG到Airflow
    如果你已经在外部代码里定义好了DAG实例,可以用Airflow自带的dag.serialize()方法把它转成符合API要求的结构化数据,然后调用POST /api/v1/dags接口发布。示例代码(用requests库):

    import requests
    from airflow.models import DAG
    from datetime import datetime
    
    # 假设这是你在外部代码里定义的DAG
    my_external_dag = DAG(
        dag_id="external_managed_dag",
        start_date=datetime(2024, 1, 1),
        schedule=None,  # 手动触发的话设为None
        catchup=False
    )
    
    # 序列化DAG为API可接受的格式
    serialized_dag = my_external_dag.serialize()
    
    # 调用Airflow Webserver API发布
    airflow_api_base = "http://your-airflow-webserver:8080/api/v1"
    auth = ("your-username", "your-password")  # 替换成你的Airflow认证信息
    
    publish_response = requests.post(
        f"{airflow_api_base}/dags",
        json=serialized_dag,
        auth=auth
    )
    publish_response.raise_for_status()
    print(f"DAG发布成功: {publish_response.json()['dag_id']}")
    
  • 第二步:触发DAG运行
    DAG发布成功后,直接调用POST /api/v1/dags/{dag_id}/dagRuns接口触发执行:

    trigger_payload = {
        "conf": {"custom_param": "hello_from_external"},  # 可选,传递运行参数
        "execution_date": datetime.utcnow().isoformat()
    }
    
    trigger_response = requests.post(
        f"{airflow_api_base}/dags/external_managed_dag/dagRuns",
        json=trigger_payload,
        auth=auth
    )
    trigger_response.raise_for_status()
    print(f"DAG触发成功,运行ID: {trigger_response.json()['dag_run_id']}")
    
如果你想继续用dag_pickle表的现有条目

如果已经通过序列化把DAG存到了dag_pickle表,需要让调度器加载这个pickle并激活DAG:

  • 首先,确保DAG在dag表中的状态是未暂停的:你可以直接在元数据库中更新dag表的is_paused字段为False,或者用API调用PATCH /api/v1/dags/{dag_id}设置:
    patch_response = requests.patch(
        f"{airflow_api_base}/dags/external_managed_dag",
        json={"is_paused": False},
        auth=auth
    )
    
  • 然后,触发调度器重新加载DAGs:可以在Airflow UI的DAG页面点击「Refresh」按钮,或者执行命令:
    airflow dags refresh
    
    也可以用API触发刷新:
    requests.post(f"{airflow_api_base}/dags/refresh", auth=auth)
    
    调度器会扫描dag_pickle表,把序列化的DAG加载到内存,之后就能正常触发了。
关键注意事项
  • 不管用哪种方式,Airflow调度器的运行环境必须能找到DAG依赖的所有资源(比如自定义Operator、导入的模块),否则会加载失败。
  • REST API方式更灵活,不需要直接操作数据库,权限控制也更清晰,优先推荐。
  • 如果是Airflow 2.x版本,确保API的认证方式配置正确(比如默认的Basic Auth,或者你自己配置的OAuth)。

内容的提问来源于stack exchange,提问作者sykik

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 10:44:20