"""DAG to release reserves for a specified statement period.""" from datetime import datetime from airflow import DAG from airflow.operators.python import PythonOperator from lib import constants from tasks.reserves_release.invoke_reserves_release_lambda import \ invoke_reserves_release_lambda_task from tasks.reserves_release.notify_failure import \ notify_failure_task from tasks.reserves_release.notify_started import \ notify_started_task from tasks.reserves_release.notify_success import \ notify_success_task dag = DAG( constants.DAG_RESERVES_RELEASE_NAME, description="Reserves are credited back to an account\'s ledger.", start_date=datetime(2019, 1, 1), default_args={'provide_context': True}, schedule_interval=None ) """ ###################################################################### ############################# OPERATORS ############################## ###################################################################### """ invoke_reserves_release_lambda_operator = PythonOperator( dag=dag, python_callable=invoke_reserves_release_lambda_task, task_id='invoke_reserves_release_processor' ) notify_failure_operator = PythonOperator( dag=dag, python_callable=notify_failure_task, task_id='notify_reserves_release_lambda_failure', trigger_rule='one_failed' ) notify_started_operator = PythonOperator( dag=dag, python_callable=notify_started_task, task_id='notify_reserves_release_lambda_started' ) notify_success_operator = PythonOperator( dag=dag, python_callable=notify_success_task, task_id='notify_reserves_release_lambda_success', trigger_rule='none_failed' ) notify_started_operator >> \ invoke_reserves_release_lambda_operator >> \ [notify_success_operator, notify_failure_operator]