"""DAG for auto-generating, validating and importing an adjustment file.""" from datetime import datetime from airflow import DAG from airflow.operators.python import PythonOperator from airflow.operators.python import ShortCircuitOperator from lib import constants from tasks.adjustment_file_generate.lambdas import invoke_generate_task from tasks.adjustment_file_generate.lambdas import invoke_import_task from tasks.adjustment_file_generate.lambdas import invoke_validate_task from tasks.adjustment_file_generate.states import failure_generate_task from tasks.adjustment_file_generate.states import failure_import_task from tasks.adjustment_file_generate.states import failure_validate_task from tasks.adjustment_file_generate.states import start_generate_task from tasks.adjustment_file_generate.states import start_import_task from tasks.adjustment_file_generate.states import start_validate_task from tasks.adjustment_file_generate.states import success_generate_task from tasks.adjustment_file_generate.states import success_import_task from tasks.adjustment_file_generate.states import success_validate_task dag = DAG( constants.DAG_ADJUSTMENT_FILE_GENERATE_NAME, description='Generate, validate and import an adjustment file', start_date=datetime(2019, 1, 1), default_args={'provide_context': True}, schedule_interval=None ) # Operators start_generate_operator = PythonOperator( dag=dag, python_callable=start_generate_task, task_id='start_generate', trigger_rule='none_failed' ) invoke_generate_operator = ShortCircuitOperator( dag=dag, python_callable=invoke_generate_task, task_id='invoke_generate', trigger_rule='none_failed' ) success_generate_operator = PythonOperator( dag=dag, python_callable=success_generate_task, task_id='success_generate', trigger_rule='none_failed' ) failure_generate_operator = ShortCircuitOperator( dag=dag, python_callable=failure_generate_task, task_id='failure_generate', trigger_rule='one_failed' ) start_validate_operator = PythonOperator( dag=dag, python_callable=start_validate_task, task_id='start_validate', trigger_rule='none_failed' ) invoke_validate_operator = PythonOperator( dag=dag, python_callable=invoke_validate_task, task_id='invoke_validate', trigger_rule='none_failed' ) success_validate_operator = PythonOperator( dag=dag, python_callable=success_validate_task, task_id='success_validate', trigger_rule='none_failed' ) failure_validate_operator = ShortCircuitOperator( dag=dag, python_callable=failure_validate_task, task_id='failure_validate', trigger_rule='one_failed' ) start_import_operator = PythonOperator( dag=dag, python_callable=start_import_task, task_id='start_import', trigger_rule='none_failed' ) invoke_import_operator = PythonOperator( dag=dag, python_callable=invoke_import_task, task_id='invoke_import', trigger_rule='none_failed' ) success_import_operator = PythonOperator( dag=dag, python_callable=success_import_task, task_id='success_import', trigger_rule='none_failed' ) failure_import_operator = ShortCircuitOperator( dag=dag, python_callable=failure_import_task, task_id='failure_import', trigger_rule='one_failed' ) # DAG Flow start_generate_operator >> \ invoke_generate_operator >> \ [success_generate_operator, failure_generate_operator] >> \ start_validate_operator >> \ invoke_validate_operator >> \ [success_validate_operator, failure_validate_operator] >> \ start_import_operator >> \ invoke_import_operator >> \ [success_import_operator, failure_import_operator]