"""Invoke lambda to send abacus payments into Payoneer system.""" from datetime import datetime from airflow import DAG from airflow.operators.python import PythonOperator from lib import constants from tasks.payoneer_payments_payout.invoke_payoneer_payments_payout_lambda import \ invoke_payoneer_payments_payout_lambda_task from tasks.payoneer_payments_payout.notify_failure import notify_failure_task from tasks.payoneer_payments_payout.notify_success import notify_success_task dag = DAG( constants.DAG_PAYONEER_PAYMENTS_PAYOUT_NAME, description='Invoke lambda to send payments to Payoneer and handle the result', start_date=datetime(2019, 1, 1), default_args={'provide_context': True}, schedule_interval=None ) """ ###################################################################### ############################# OPERATORS ############################## ###################################################################### """ invoke_payoneer_payments_payout_lambda_operator = PythonOperator( task_id='invoke_payoneer_payments_payout_lambda', python_callable=invoke_payoneer_payments_payout_lambda_task, dag=dag ) notify_failure_operator = PythonOperator( dag=dag, python_callable=notify_failure_task, task_id='notify_failure', op_kwargs={'action_name': constants.PAYMENT_GROUP_PAYMENT_ACTIONS.SEND_PAYMENTS}, trigger_rule='one_failed' ) notify_success_operator = PythonOperator( dag=dag, python_callable=notify_success_task, op_kwargs={'action_name': constants.PAYMENT_GROUP_PAYMENT_ACTIONS.SEND_PAYMENTS}, task_id='notify_success' ) invoke_payoneer_payments_payout_lambda_operator >> \ [notify_success_operator, notify_failure_operator]