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

Airflow 2.0中能否用TaskFlow API的@task装饰器实现任务回调?

关于Airflow 2.0 TaskFlow API回调机制的说明

Airflow 2.0的TaskFlow API完全支持任务成功/失败的回调机制,但你给出的代码写法存在两处问题,调整后即可正常使用:

  • 回调函数需与任务函数分离
    不能把被@task装饰的任务函数同时设为回调函数,两者逻辑要分开:回调是任务执行完成后触发的独立逻辑,任务函数是核心业务逻辑。

  • 回调参数需传入函数对象而非字符串
    配置on_success_callback时,要直接传函数名(不带引号),不能传字符串形式的函数名。

正确的示例代码

from airflow.decorators import task
from airflow import DAG
import pendulum

# 定义任务成功后的回调函数
def task_success_alert(context):
    print(f"任务执行成功,任务ID:{context['task_instance'].task_id}")

# 定义任务失败后的回调函数
def task_failure_alert(context):
    print(f"任务执行失败,任务ID:{context['task_instance'].task_id}")

with DAG(
    dag_id="taskflow_callback_demo",
    start_date=pendulum.datetime(2021, 1, 1, tz="UTC"),
    schedule=None,
    catchup=False
):
    # 使用@task装饰器指定回调参数
    @task(
        task_id="business_task",
        on_success_callback=task_success_alert,
        on_failure_callback=task_failure_alert
    )
    def business_logic_task():
        # 这里写实际的任务业务逻辑
        print("执行核心业务任务")
    
    # 触发任务
    business_logic_task()

额外说明

  • 回调支持传入函数列表,比如on_success_callback=[func_a, func_b],列表中的函数会按顺序依次执行
  • DAG层面的全局回调(on_success_callback、on_failure_callback)在TaskFlow API中同样适用,直接在DAG的定义参数里指定即可

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 02:05:04