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
相关产品推荐
相关产品推荐

