Trigger DAG run from a task in another DAG
December 21, 2018
When writing airflow jobs sometimes it’s necessary to trigger a DAG run from a task in another DAG. Airfow provides TriggerDagRunOperator for this purpose which is great. A caveat with TriggerDagRunOperator is that it will add a DAG run object whether the target DAG is paused or not. This will cause the target DAG to start executing as soon as it’s enabled which is not always desirable.
The following code example shows how you can avoid that by first checking whether a DAG is paused before creating the DAG run.
from airflow.settings import Session
from airflow.models import DagModel
from airflow.operators.dagrun_operator import TriggerDagRunOperator
import logging
from etl import utils
def get_dag_with_id(dag_id):
session = Session()
try:
qry = session.query(DagModel).filter(DagModel.dag_id == dag_id)
dag = qry.first()
finally:
session.close()
return dag
def create_trigger_dag_run_operator(dag, args, target_dag_id, task_id, trigger_rule='all_success'):
def trigger_dag_run(context, dag_run_obj):
"""This function decides whether or not to Trigger the remote DAG"""
target_dag = get_dag_with_id(target_dag_id)
if target_dag.is_paused:
logging.info('target dag {} is paused. will not create a dag run object'.format(target_dag))
return None
else:
dag_run_timestamp = utils.current_date_stamp_dag_run_format()
dag_run_id = 'triggered__{}'.format(dag_run_timestamp)
logging.info('trigger target dag {0} with run id {1}'.format(target_dag_id, dag_run_id))
if context['params']['do_trigger']:
dag_run_obj.run_id = dag_run_id
dag_run_obj.payload = {
'triggered_from_dag': context['params']['triggered_from_dag'],
'trigger_timestamp': dag_run_timestamp
}
return dag_run_obj
return TriggerDagRunOperator(
task_id=task_id,
trigger_dag_id=target_dag_id,
python_callable=trigger_dag_run,
params={
'do_trigger': True,
'triggered_from_dag': dag.dag_id
},
trigger_rule=trigger_rule,
dag=dag,
default_args=args
)Written by Francois Fernando, a software craftsman and tinkerer.