"""DAG for calculating NR sales against PNR contracts.""" from datetime import datetime from airflow import DAG from airflow.operators.dummy import DummyOperator from airflow.operators.python import PythonOperator from lib import constants from tasks.accounting_run_calculate_nr.invoke_ledger_accounting_run_balance_lambda \ import invoke_ledger_accounting_run_balance_lambda_task from tasks.accounting_run_calculate_nr.snowflake_calculate_totals \ import calculate_run_totals_task from tasks.accounting_run_calculate_nr.snowflake_flatten_contract_json \ import insert_flat_contract_task from tasks.accounting_run_calculate_nr.snowflake_match_contract_terms \ import match_contract_terms_to_sales_task from tasks.accounting_run_calculate_nr.snowflake_snapshot_contract_term_schedules \ import snapshot_contract_term_schedules_nr_task from tasks.accounting_run_calculate_nr.snowflake_snapshot_contracts \ import insert_snapshot_contracts_nr_task from tasks.accounting_run_calculate_nr.snowflake_snapshot_flat_contract_term_conditions\ import snapshot_flat_contract_term_conditions_task from tasks.accounting_run_calculate_nr.update_run_status_complete import status_complete from tasks.accounting_run_calculate_nr.update_run_status_error import status_error from tasks.accounting_run_calculate_nr.update_run_status_started import status_started dag = DAG( constants.DAG_NR_CALC_EVENT_NAME, description='Calculate NR revenue based on contract terms', start_date=datetime(2023, 10, 1), default_args={'provide_context': True}, schedule_interval=None ) """ ###################################################################### ############################# OPERATORS ############################## ###################################################################### """ dummy_join_flatten_term_operator = DummyOperator( dag=dag, task_id='insert_flat_contract_data' ) invoke_ledger_accounting_run_balance_lambda_operator = PythonOperator( dag=dag, python_callable=invoke_ledger_accounting_run_balance_lambda_task, task_id='invoke_ledger_accounting_run_balance_lambda', trigger_rule='none_failed' ) snowflake_calculate_totals_operator = PythonOperator( dag=dag, python_callable=calculate_run_totals_task, task_id='snowflake_calculate_run_totals' ) snowflake_flatten_contract_json_operator = PythonOperator( dag=dag, op_kwargs={'is_contributor_only': False}, python_callable=insert_flat_contract_task, task_id='snowflake_insert_flat_contract_data' ) snowflake_flatten_contrib_only_contract_json_operator = PythonOperator( dag=dag, op_kwargs={'is_contributor_only': True}, python_callable=insert_flat_contract_task, task_id='snowflake_insert_flat_contributor_only_contract_data' ) snowflake_match_contract_terms_operator = PythonOperator( dag=dag, op_kwargs={'is_contributor_only': False}, python_callable=match_contract_terms_to_sales_task, task_id='snowflake_match_contract_terms_to_sales' ) snowflake_match_contrib_only_contract_terms_operator = PythonOperator( dag=dag, op_kwargs={'is_contributor_only': True}, python_callable=match_contract_terms_to_sales_task, task_id='snowflake_match_contributor_only_contract_terms_to_sales' ) snowflake_snapshot_contracts_operator = PythonOperator( dag=dag, python_callable=insert_snapshot_contracts_nr_task, task_id='snapshot_contracts_to_snowflake' ) snowflake_snapshot_flat_contract_term_condition_countries_operator = PythonOperator( dag=dag, op_kwargs={'condition_type': constants.CONTRACT_TERM_CONDITIONS.COUNTRIES}, python_callable=snapshot_flat_contract_term_conditions_task, task_id='snowflake_snapshot_flat_countries' ) snowflake_snapshot_flat_contract_term_condition_stores_operator = PythonOperator( dag=dag, op_kwargs={'condition_type': constants.CONTRACT_TERM_CONDITIONS.STORES}, python_callable=snapshot_flat_contract_term_conditions_task, task_id='snowflake_snapshot_flat_stores' ) snowflake_snapshot_flat_contract_term_condition_transaction_types_operator = \ PythonOperator( dag=dag, op_kwargs={ 'condition_type': constants.CONTRACT_TERM_CONDITIONS.TRANSACTION_TYPES }, python_callable=snapshot_flat_contract_term_conditions_task, task_id='snowflake_snapshot_flat_transaction_types' ) snowflake_snapshot_contract_term_schedules_operator = PythonOperator( dag=dag, python_callable=snapshot_contract_term_schedules_nr_task, task_id='snapshot_contract_term_schedules_to_snowflake' ) status_complete_operator = PythonOperator( dag=dag, python_callable=status_complete, task_id='update_status_complete', trigger_rule='none_failed' ) status_error_operator = PythonOperator( dag=dag, python_callable=status_error, task_id='status_error', trigger_rule='one_failed', ) status_started_operator = PythonOperator( dag=dag, python_callable=status_started, task_id='status_started' ) """ ###################################################################### ################################ DAG ################################# ###################################################################### """ status_started_operator >> \ snowflake_snapshot_contracts_operator >> \ snowflake_snapshot_contract_term_schedules_operator >> [ snowflake_snapshot_flat_contract_term_condition_countries_operator, snowflake_snapshot_flat_contract_term_condition_stores_operator, snowflake_snapshot_flat_contract_term_condition_transaction_types_operator ] >> dummy_join_flatten_term_operator >> [ snowflake_flatten_contract_json_operator, snowflake_flatten_contrib_only_contract_json_operator ] >> \ snowflake_match_contract_terms_operator >> \ snowflake_match_contrib_only_contract_terms_operator >> \ snowflake_calculate_totals_operator >> \ invoke_ledger_accounting_run_balance_lambda_operator >> \ status_complete_operator >> \ status_error_operator