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

