【问题标题】:How is exception passed to on_failure_callback?异常如何传递给 on_failure_callback?
【发布时间】:2021-10-24 10:27:48
【问题描述】:

我想将异常传递给 on_failure_callback 以检查错误是什么。例如,如果它在某个 DAG 中包含“存在重复项”,则该函数不会执行任何操作。否则,它将发送一封电子邮件。

但是,我看不到异常的格式。我在 Docker 中使用 Airflow 2.1.2 并作为我的 dag 定义如下:

with DAG(process_name,
         default_args=default_args,
         schedule_interval='@daily',
         max_active_runs=1,
         tags=['import', 'es'],
         on_failure_callback=known_error_dag
         ) as dag:
    operators

已经尝试了下一个解决方案:

def known_error_dag(context):
    #  1
    ti = context['ti']
    ti.xcom_push(key='exception', value=context['exception'])

    #  2
    print(context['exception'])
    
    #  3
    logging.info(context['exception'])

我在 UI 和 docker 日志中都看不到异常。此外,它不会出现在 XCOM 中。

这个问题的答案不清楚我想要什么是可能的:Get Exception details on Airflow on_failure_callback context

然而,天文学家课程指出这确实是可能的。 https://academy.astronomer.io/astronomer-certification-apache-airflow-dag-authoring-preparation

【问题讨论】:

    标签: airflow


    【解决方案1】:

    您可以在 DAG 和任务级别上定义 on_failure_callback。异常仅传递给任务级别的失败回调,因此在您的操作员上或通过 DAG 上的 default_args 配置回调到所有操作员:

    with DAG(
        process_name,
        default_args=default_args,
        schedule_interval='@daily',
        max_active_runs=1,
        tags=['import', 'es'],
        default_args={
            "on_failure_callback": known_error_dag,
        },
    ) as dag:
    

    在 DAG 级别定义的 on_failure_callback 也将采用上下文变量,其中包括一个键“reason”,但仅在 DAG 运行失败的情况下声明“task_failure”,因此不是很有用在大多数情况下。

    【讨论】:

    • 我已经尝试过,但仍然无法正常工作。如果我进入失败任务的任务实例详细信息菜单,on_failure_callback 部分确实有我定义的功能,但我看不到任何日志/打印并且 xcom 仍然不存在
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-03-15
    • 1970-01-01
    • 1970-01-01
    • 2013-01-12
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多