"""DAG for getting eligible sales from snowflake and creating a sales file.""" from datetime import datetime from airflow import DAG from airflow.operators.python import PythonOperator from lib.constants import ABACUS_STATE_STATUSES from lib.constants import DAG_SALES_GET_ELIGIBLE_NAME from lib.utils.slack import slack_failure_callback from tasks.sales_get_eligible.snowflake_get_eligible_sales import get_eligible_sales from tasks.sales_get_eligible.update_abacus_state import update_abacus_state_task from tasks.sales_get_eligible.update_sales_file_metadata import update_sales_metadata dag = DAG( DAG_SALES_GET_ELIGIBLE_NAME, default_args={ 'provide_context': True, 'on_failure_callback': slack_failure_callback, }, description='Query snowflake for eligible sales and create a sales file record.', schedule_interval=None, start_date=datetime(2022, 1, 1) ) """ ###################################################################### ############################# OPERATORS ############################## ###################################################################### """ get_eligible_sales_operator = PythonOperator( dag=dag, python_callable=get_eligible_sales, task_id='get_eligible_sales' ) notify_failure_operator = PythonOperator( dag=dag, op_kwargs={'action_status': ABACUS_STATE_STATUSES.ERROR}, python_callable=update_abacus_state_task, task_id='notify_failure', trigger_rule='one_failed' ) notify_started_operator = PythonOperator( dag=dag, op_kwargs={'action_status': ABACUS_STATE_STATUSES.RUNNING}, python_callable=update_abacus_state_task, task_id='notify_started' ) notify_success_operator = PythonOperator( dag=dag, op_kwargs={'action_status': ABACUS_STATE_STATUSES.COMPLETE}, python_callable=update_abacus_state_task, task_id='notify_success', trigger_rule='none_failed' ) update_sales_file_metadata_operator = PythonOperator( dag=dag, python_callable=update_sales_metadata, task_id='update_sales_file_metadata' ) """ ###################################################################### ################################ DAG ################################# ###################################################################### """ notify_started_operator >> \ get_eligible_sales_operator >> \ update_sales_file_metadata_operator >> \ notify_success_operator >> \ notify_failure_operator