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

Airflow中如何实现任务失败分支执行并获取错误日志?

在Airflow中实现带错误日志捕获的分支逻辑(If-Else风格)

嘿,我懂你要实现的需求——在Airflow里搞个类似if-else的分支:task_1成功就跑task_3,失败就触发task_2,还得能拿到task_1的错误日志对吧?之前用SSHHOOK觉得麻烦,那我给你分享个更直接的方案,完全用Airflow内置的功能就能搞定!

核心思路

我们可以用BranchPythonOperator来做分支判断,同时借助Airflow的TaskInstance类直接获取task_1的执行状态和错误日志,不需要额外的SSH操作。整个流程是:

  1. 执行task_1
  2. 分支任务判断task_1的状态:
    • 成功 → 触发task_3
    • 失败 → 捕获task_1的错误日志,传递给task_2后触发它

完整代码示例

from airflow.models import TaskInstance
from airflow.operators.python import PythonOperator, BranchPythonOperator
from airflow.utils.state import State
from datetime import datetime

def determine_next_task(**context):
    # 从上下文获取task_1的TaskInstance对象
    ti_task1 = TaskInstance(
        task_id='task_1',
        dag_id=context['dag'].dag_id,
        execution_date=context['execution_date']
    )
    # 从数据库刷新最新的任务状态
    ti_task1.refresh_from_db()
    
    if ti_task1.state == State.SUCCESS:
        # task_1成功,返回task_3的task_id
        return 'task_3'
    else:
        # 捕获task_1的错误日志
        error_log = ti_task1.get_log()
        # 把日志通过XCom传递给后续任务
        context['ti'].xcom_push(key='task_1_error_log', value=error_log)
        # 返回task_2的task_id
        return 'task_2'

with DAG(
    'airflow_branch_with_error_log',
    catchup=False,
    default_args={
        'owner': 'abc',
        'start_date': datetime(2018, 4, 17),
        'schedule_interval': None,
        'depends_on_past': False,
    },
) as dag:
    # 模拟task_1:这里故意写个除零错误模拟失败,成功的话改成lambda: print("Task1 finished successfully")
    task_1 = PythonOperator(
        task_id='task_1',
        python_callable=lambda: 1/0,
        provide_context=True
    )

    # 分支判断任务
    branch_decision = BranchPythonOperator(
        task_id='branch_decision',
        python_callable=determine_next_task,
        provide_context=True
    )

    # task_2:接收并打印task_1的错误日志
    task_2 = PythonOperator(
        task_id='task_2',
        python_callable=lambda **context: print(f"Task1 failed! Error log:\n{context['ti'].xcom_pull(key='task_1_error_log', task_ids='branch_decision')}"),
        provide_context=True
    )

    # task_3:task_1成功后执行的任务
    task_3 = PythonOperator(
        task_id='task_3',
        python_callable=lambda: print("Task1 succeeded! Running task_3 now.")
    )

    # 设置任务依赖关系
    task_1 >> branch_decision >> [task_2, task_3]

关键细节说明

  • provide_context=True:让PythonOperator能获取Airflow的上下文信息(比如execution_date、dag ID等),这是获取TaskInstance和传递XCom的前提。
  • TaskInstance.refresh_from_db():确保我们拿到的是task_1最新的执行状态,避免因为缓存导致判断错误。
  • ti.get_log():直接从Airflow的日志存储中获取task_1的完整日志,不需要额外的SSH连接。
  • XCom传递日志:如果日志内容较大,建议把日志写入外部存储(比如S3、数据库),然后只传递存储路径给task_2,避免XCom的大小限制问题。

注意事项

如果你的Airflow用的是分布式Executor(比如CeleryExecutor),要确保日志存储是全局可访问的(比如用统一的日志系统如Elasticsearch、S3),这样ti.get_log()才能正常获取到日志内容。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 04:27:08