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

Apache Airflow 如何在每次DAG运行成功后调用外部许可服务器

实现方案

Airflow原生支持多种方式实现DAG运行成功后触发指定逻辑,你可以根据实际场景选择对应方案:

方案1:单DAG配置成功回调(最常用,适合单个/少数DAG需求)

直接给DAG对象配置on_success_callback参数即可,回调函数里写调用外部许可服务器的逻辑就行。
示例代码如下:

from airflow import DAG
from airflow.operators.bash import BashOperator
from datetime import datetime
import requests

# 自定义回调函数:调用外部许可服务器
def call_license_server(context):
    # 可从context中提取DAG运行相关参数,按需传给许可服务器
    dag_run_id = context['dag_run'].run_id
    exec_date = context['execution_date']
    try:
        # 替换为实际的许可服务器接口地址和请求参数
        resp = requests.post(
            "http://your-license-server/api/check",
            json={
                "dag_id": context['dag'].dag_id,
                "run_id": dag_run_id,
                "execution_date": exec_date.isoformat()
            },
            timeout=10
        )
        resp.raise_for_status()
    except Exception as e:
        # 可按需添加异常处理,比如日志打印、告警推送
        print(f"调用许可服务器失败: {str(e)}")
        raise

# 定义DAG时指定on_success_callback
with DAG(
    dag_id="your_dag_name",
    start_date=datetime(2024, 1, 1),
    schedule_interval="@daily",
    on_success_callback=call_license_server, # 核心配置项
    catchup=False
) as dag:
    # 你的正常DAG任务
    task1 = BashOperator(
        task_id="task1",
        bash_command="echo 'task1 running'"
    )

    task2 = BashOperator(
        task_id="task2",
        bash_command="echo 'task2 running'"
    )

    task1 >> task2

该方案灵活度高,每个DAG可单独配置不同的回调逻辑,不需要修改全局配置。

方案2:全局配置所有DAG默认触发(适合所有DAG都需要调用许可服务器的场景)

如果希望所有DAG运行成功后自动调用,不需要每个DAG单独加配置,可以用Airflow的policy机制:

  1. 找到Airflow配置目录下的airflow_local_settings.py文件,不存在则新建
  2. 添加如下代码:
from airflow.models import DAG
import requests

def call_license_server(context):
    # 和方案1的回调函数逻辑一致
    dag_run_id = context['dag_run'].run_id
    resp = requests.post(
        "http://your-license-server/api/check",
        json={"dag_id": context['dag'].dag_id, "run_id": dag_run_id},
        timeout=10
    )
    resp.raise_for_status()

def dag_policy(dag: DAG):
    # 给所有DAG默认添加成功回调,可自行选择覆盖还是追加原有回调
    if not dag.on_success_callback:
        dag.on_success_callback = call_license_server
    else:
        # 保留原有回调同时新增许可调用的写法
        original_callback = dag.on_success_callback
        def combined_callback(context):
            original_callback(context)
            call_license_server(context)
        dag.on_success_callback = combined_callback
  1. 重启Airflow的scheduler服务即可生效。

方案3:末尾专用触发任务(适合需要把调用逻辑作为可见任务展示在DAG流程里的场景)

如果希望调用许可服务器的操作能在DAG运行图中直观看到,也可以专门写一个任务放在DAG末尾,设置触发规则为所有上游成功:

from airflow import DAG
from airflow.operators.bash import BashOperator
from airflow.operators.python import PythonOperator
from datetime import datetime
import requests

def call_license_server(**kwargs):
    dag_run = kwargs['dag_run']
    resp = requests.post(
        "http://your-license-server/api/check",
        json={"dag_id": kwargs['dag'].dag_id, "run_id": dag_run.run_id},
        timeout=10
    )
    resp.raise_for_status()

with DAG(
    dag_id="your_dag_name",
    start_date=datetime(2024, 1, 1),
    schedule_interval="@daily",
    catchup=False
) as dag:
    task1 = BashOperator(task_id="task1", bash_command="echo 'task1'")
    task2 = BashOperator(task_id="task2", bash_command="echo 'task2'")
    # 许可调用专用任务
    license_task = PythonOperator(
        task_id="call_license_server",
        python_callable=call_license_server,
        trigger_rule="all_success" # 仅所有上游任务成功时才执行
    )

    # 把所有原有任务指向许可任务
    [task1, task2] >> license_task
注意事项
  • 如果许可服务器接口响应较慢,建议给回调函数加合理的超时配置,避免阻塞Airflow scheduler
  • 如果许可调用失败需要重试,方案3的PythonOperator可直接配置retries参数实现重试,方案1和2的回调需要自行在函数中添加重试逻辑
  • 许可服务器的鉴权token等敏感信息不要硬编码在代码里,建议用Airflow的Variable或者Connection存储

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 03:54:02