如何调用DAG任务,跨DAG复用已有任务创建新DAG
实现方案(Airflow 场景)
以下两种方案均无需重写原有任务的业务逻辑代码即可实现需求:
方案1:直接复用原有任务函数(最简便)
如果task A、E、C的函数是单独定义、未和Dag1、Dag2强绑定的,直接导入函数到Dag3的定义文件中实例化即可:
- 步骤1:在Dag3文件头部导入三个任务对应的函数,示例:
from dag1_def import task_a_func、from dag2_def import task_e_func、from dag1_def import task_c_func - 步骤2:在Dag3的上下文里重新实例化三个任务,直接设置依赖关系即可,示例代码如下:
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime with DAG( dag_id='Dag3', start_date=datetime(2024,1,1), schedule_interval=None, catchup=False ) as dag: task_a = PythonOperator( task_id='task_A', python_callable=task_a_func, # 直接复用原有task A的op_args、op_kwargs等配置即可 ) task_e = PythonOperator( task_id='task_E', python_callable=task_e_func, # 直接复用原有task E的配置即可 ) task_c = PythonOperator( task_id='task_C', python_callable=task_c_func, # 直接复用原有task C的配置即可 ) # 设置执行顺序 task_a >> task_e >> task_c
该方案下Dag3的三个任务是独立实例,和Dag1、Dag2的原有任务运行互不干扰。
方案2:依赖原有DAG的任务运行(触发原有任务实例)
如果需要Dag3实际触发的是Dag1、Dag2中原有的任务实例,而非新建独立任务,可以用TriggerDagRunOperator + ExternalTaskSensor组合实现:
- 步骤1:在Dag1、Dag2中配置支持按参数触发指定任务的逻辑
- 步骤2:在Dag3中依次触发对应任务、监听运行状态,确认成功后再触发下一个任务,核心示例代码如下:
from airflow import DAG from airflow.operators.trigger_dagrun import TriggerDagRunOperator from airflow.sensors.external_task import ExternalTaskSensor from datetime import datetime with DAG( dag_id='Dag3', start_date=datetime(2024,1,1), schedule_interval=None, catchup=False ) as dag: # 触发并等待Dag1的task A运行完成 trigger_a = TriggerDagRunOperator( task_id='trigger_task_a', trigger_dag_id='Dag1', execution_date="{{ execution_date }}", conf={"run_only_task": "task_A"} ) wait_a = ExternalTaskSensor( task_id='wait_a_finish', external_dag_id='Dag1', external_task_id='task_A', execution_date_fn=lambda dt: dt, timeout=3600 ) # 触发并等待Dag2的task E运行完成 trigger_e = TriggerDagRunOperator( task_id='trigger_task_e', trigger_dag_id='Dag2', execution_date="{{ execution_date }}", conf={"run_only_task": "task_E"} ) wait_e = ExternalTaskSensor( task_id='wait_e_finish', external_dag_id='Dag2', external_task_id='task_E', execution_date_fn=lambda dt: dt, timeout=3600 ) # 触发并等待Dag1的task C运行完成 trigger_c = TriggerDagRunOperator( task_id='trigger_task_c', trigger_dag_id='Dag1', execution_date="{{ execution_date }}", conf={"run_only_task": "task_C"} ) wait_c = ExternalTaskSensor( task_id='wait_c_finish', external_dag_id='Dag1', external_task_id='task_C', execution_date_fn=lambda dt: dt, timeout=3600 ) # 设置执行顺序 trigger_a >> wait_a >> trigger_e >> wait_e >> trigger_c >> wait_c
该方案需要Dag1、Dag2的调度规则和Dag3对齐,避免sensor监听不到对应任务实例。
内容的提问来源于stack exchange,提问作者The Doctor
相关产品推荐
相关产品推荐

