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

Airflow任务中如何获取其他任务实例信息及任务组完成状态?

解决方案

一、修正代码基础错误

首先修复call_task_2重复的task_id问题,改为唯一值call_task_2,否则DAG无法正常加载运行。

二、在call_task_2中获取call_task_1的实例信息

要获取其他任务的实例状态、执行日期等信息,需通过Airflow的TaskInstance模型从元数据库查询。在get_task_2_log函数中,利用当前任务上下文定位到call_task_1的实例:

核心实现步骤:

  • 导入airflow.models.TaskInstance模块
  • 通过kwargs中的dag_run获取当前DAG的执行日期
  • 使用完整任务ID(TaskGroup内任务格式为{组名}.{任务ID})查询目标任务实例,提取所需信息

三、判断TaskGroup内所有任务完成后推进后续任务组

TaskGroup本身是一个逻辑单元,默认触发规则为all_success,只需将后续任务组与当前TaskGroup建立依赖关系,即可实现“当前组所有任务完成后才执行下一组”的逻辑。若需兼容任务失败的场景,可给TaskGroup设置trigger_rule="all_done"参数。

完整修正后的代码

import os, json, time, airflow, requests
from airflow import DAG
from datetime import datetime, timedelta, timezone
from airflow.configuration import conf
from airflow.models import Variable, TaskInstance
from airflow.utils.task_group import TaskGroup
from airflow.operators.python_operator import PythonOperator

# 任务执行函数
def get_task_1_log(**kwargs):
    task_instance = kwargs['task_instance']
    print(f"call_task_1 task_id: {task_instance.task_id}")
    print(f"call_task_1 dag_id: {task_instance.dag_id}")
    print(f"call_task_1 execution_date: {task_instance.execution_date}")

def get_task_2_log(**kwargs):
    task_instance = kwargs['task_instance']
    dag_run = kwargs['dag_run']
    print(f"call_task_2 task_id: {task_instance.task_id}")
    print(f"call_task_2 dag_id: {task_instance.dag_id}")
    print(f"call_task_2 execution_date: {task_instance.execution_date}")

    # 获取call_task_1的实例信息
    target_task_full_id = "my_group.call_task_1"
    ti = TaskInstance.find(
        dag_id=dag_run.dag_id,
        task_id=target_task_full_id,
        execution_date=dag_run.execution_date
    )[0]
    print(f"call_task_1 状态: {ti.state}")
    print(f"call_task_1 执行日期: {ti.execution_date}")
    print(f"call_task_1 开始时间: {ti.start_date}")
    print(f"call_task_1 结束时间: {ti.end_date}")

# 默认参数初始化
default_args = {
    'start_date': datetime(2024, 1, 27),
    'retries': 1,
    'retry_delay': timedelta(minutes=5)
}

# DAG定义
with DAG("Get_Task_Logs",
        default_args=default_args,
        description="Get_Task_Logs",
        schedule_interval="07 06 * * *",
        start_date=None,
        ) as dag:
    
    # 第一个任务组
    with TaskGroup("my_group", tooltip="my_group") as my_group:
        call_task_1 = PythonOperator(
            task_id="call_task_1",
            python_callable=get_task_1_log,
            trigger_rule='one_success'
        )
        
        call_task_2 = PythonOperator(
            task_id="call_task_2",  # 修正重复的task_id
            python_callable=get_task_2_log,
            trigger_rule='one_success'
        )
        
        call_task_1 >> call_task_2 
    
    # 示例后续任务组
    with TaskGroup("next_group", tooltip="后续任务组") as next_group:
        def follow_up_task(**kwargs):
            print("执行后续任务组的任务")
        
        follow_up = PythonOperator(
            task_id="follow_up",
            python_callable=follow_up_task
        )
    
    # 依赖设置:my_group所有任务完成后执行next_group
    my_group >> next_group

关键注意点:

  • TaskGroup内的任务必须使用{组名}.{任务ID}的完整ID进行查询,否则无法定位到目标任务实例
  • Airflow 2.x中TaskInstance.state返回的状态枚举值包括success、failed、running等
  • 若需忽略任务失败状态,只要组内所有任务执行完毕就推进后续任务,可给TaskGroup添加参数trigger_rule="all_done"

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 05:22:34