THIS TOO SHALL PASS

Trigger DAG run from a task in another DAG

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.