"""DAG for approving sales file, moving data in snowflake, and exporting sales to S3.""" 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_APPROVE_NAME from lib.utils.slack import slack_failure_callback from tasks.sales_approve.snowflake_copy_approved_sales import copy_approved_sales from tasks.sales_approve.snowflake_export_sales import export_sales_from_snowflake from tasks.sales_approve.trigger_approve_sales_files import trigger_approve_sales_files_task # noqa: E501 from tasks.sales_approve.update_abacus_state import update_abacus_state_task dag = DAG( DAG_SALES_APPROVE_NAME, default_args={ 'provide_context': True, 'on_failure_callback': slack_failure_callback, }, description='Move data in snowflake and create parquet file of sales on S3.', schedule_interval=None, start_date=datetime(2022, 1, 1) ) """ ###################################################################### ############################# OPERATORS ############################## ###################################################################### """ 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' ) snowflake_copy_approved_sales_operator = PythonOperator( dag=dag, python_callable=copy_approved_sales, task_id='copy_approved_sales_in_snowflake' ) snowflake_export_sales_to_s3_operator = PythonOperator( dag=dag, python_callable=export_sales_from_snowflake, task_id='export_sales_from_snowflake_to_s3' ) trigger_approve_sales_files_operator = PythonOperator( dag=dag, python_callable=trigger_approve_sales_files_task, task_id='trigger_approve_sales_files', ) """ ###################################################################### ################################ DAG ################################# ###################################################################### """ notify_started_operator >> \ snowflake_copy_approved_sales_operator >> \ snowflake_export_sales_to_s3_operator >> \ notify_success_operator >> \ trigger_approve_sales_files_operator >> \ notify_failure_operator