"""DAG for running VAT's calculation flow.""" from datetime import datetime from airflow import DAG from airflow.operators.python import PythonOperator from lib import constants from tasks.accounting_period_calculate_vat.create_accounting_period_report import \ create_accounting_period_report_task from tasks.accounting_period_calculate_vat.create_vat_applied_report import \ create_vat_applied_report_task from tasks.accounting_period_calculate_vat.create_vat_exempt_report import \ create_vat_exempt_report_task from tasks.accounting_period_calculate_vat.invoke_vat_calculate_lambda import \ invoke_vat_calculate_lambda_task from tasks.accounting_period_calculate_vat.invoke_vat_exempt_lambda import \ invoke_vat_exempt_lambda_task from tasks.accounting_period_calculate_vat.notify_failure import \ notify_failure_task from tasks.accounting_period_calculate_vat.notify_success import \ notify_success_task from tasks.accounting_period_calculate_vat.snapshot_account_payment_terms import \ snapshot_account_payment_terms_task from tasks.accounting_period_calculate_vat.snapshot_account_tax_info import \ snapshot_account_tax_info_task from tasks.accounting_period_calculate_vat.snapshot_run_summaries import \ snapshot_run_summaries_task dag = DAG( constants.DAG_CALC_VAT_EVENT_NAME, description='Calculate Value Added Tax (VAT) Rates', start_date=datetime(2019, 1, 1), default_args={'provide_context': True}, schedule_interval=None ) """ ###################################################################### ############################# OPERATORS ############################## ###################################################################### """ snapshot_run_summaries_operator = PythonOperator( task_id='snapshot_run_summaries', python_callable=snapshot_run_summaries_task, dag=dag ) snapshot_account_tax_info_operator = PythonOperator( task_id='snapshot_account_tax_info', python_callable=snapshot_account_tax_info_task, dag=dag ) snapshot_account_payment_terms_operator = PythonOperator( task_id='snapshot_account_payment_terms', python_callable=snapshot_account_payment_terms_task, dag=dag ) invoke_vat_calculate_lambda_operator = PythonOperator( task_id='invoke_vat_calculate_lambda', python_callable=invoke_vat_calculate_lambda_task, dag=dag ) invoke_vat_exempt_lambda_operator = PythonOperator( task_id='invoke_vat_exempt_lambda', python_callable=invoke_vat_exempt_lambda_task, dag=dag ) create_vat_exempt_report_operator = PythonOperator( task_id='create_vat_exempt_report', python_callable=create_vat_exempt_report_task, dag=dag, trigger_rule='none_failed' ) create_vat_applied_report_operator = PythonOperator( task_id='create_vat_applied_report', python_callable=create_vat_applied_report_task, dag=dag, trigger_rule='none_failed' ) create_accounting_period_report_operator = PythonOperator( task_id='create_accounting_period_report', python_callable=create_accounting_period_report_task, dag=dag, trigger_rule='none_failed' ) notify_failure_operator = PythonOperator( task_id='notify_failure', python_callable=notify_failure_task, dag=dag, trigger_rule='one_failed' ) notify_success_operator = PythonOperator( task_id='notify_success', python_callable=notify_success_task, dag=dag, trigger_rule='none_failed' ) """ ###################################################################### ################################ DAG ################################# ###################################################################### """ invoke_lambdas = [ invoke_vat_calculate_lambda_operator, invoke_vat_exempt_lambda_operator ] snapshot_run_summaries_operator >> snapshot_account_tax_info_operator >> \ snapshot_account_payment_terms_operator >> invoke_lambdas >> \ create_vat_exempt_report_operator >> create_vat_applied_report_operator >> \ create_accounting_period_report_operator >> \ notify_success_operator >> notify_failure_operator