from airflow.models import Variable import dates as dt import slack_util from airflow import DAG from airflow.operators.dagrun_operator import TriggerDagRunOperator from airflow.operators.dummy_operator import DummyOperator from airflow.operators.python_operator import BranchPythonOperator from airflow.operators.python_operator import PythonOperator dag_name = 'apple_past_dwnlds_check_v1_0' # Dag level settings default_args = { 'owner': 'airflow', 'depends_on_past': True, 'start_date': dt.get_start_date_for_dag(), 'on_failure_callback':slack_util.dag_failure } #schedule_interval = '* * 1 * *' schedule_interval = None #schedule_interval = None dag = DAG(dag_name, default_args=default_args, max_active_runs=1, catchup=False, \ schedule_interval=schedule_interval) def is_there_dwnld_difference(date): return False def check_past_dwnlds(**kwargs): how_many_days_to_go_back = Variable.get('apple_how_many_days_to_go_back') date_range = get_date_range(int(how_many_days_to_go_back)) proceed_call_manual_dag = False for date in date_range: if is_there_dwnld_difference(date): proceed_call_manual_dag = True Variable.set("apple_manual_date_to_run", date) break if proceed_call_manual_dag: kwargs['ti'].xcom_push(key='next_task_id', value='run_apple_manual_dag') else: kwargs['ti'].xcom_push(key='next_task_id', value='no_past_dwnlds') check_past_dwnlds = PythonOperator( task_id='check_past_dwnlds', python_callable=check_past_dwnlds, provide_context=True, dag=dag ) #This is the branching function whether to proceed for downloads or not def branch_for_call_manual_dag(**kwargs): return kwargs['ti'].xcom_pull(key='next_task_id') #This is the branching python operator whether to proceed for downloads or not fork_dwnlds = BranchPythonOperator( task_id='branch_for_call_manual_dag', python_callable=branch_for_call_manual_dag, provide_context=True, dag=dag ) no_past_dwnlds = DummyOperator( task_id='no_past_dwnlds', dag=dag ) dag_run_task = TriggerDagRunOperator(task_id='run_apple_manual_dag', trigger_dag_id='apple_manual_etl_java_v1_0', dag=dag) check_past_dwnlds >> fork_dwnlds fork_dwnlds >> dag_run_task fork_dwnlds >> no_past_dwnlds