"""DAG for ingesting sales data from StatementDB to Snowflake.""" from datetime import datetime from airflow import DAG from airflow.operators.python import PythonOperator from airflow.operators.python import ShortCircuitOperator from lib import constants from lib.utils.slack import slack_failure_callback from tasks.sales_ingest.determine_sales_ingest_type import \ determine_sales_ingest_type_task from tasks.sales_ingest.run_extract_sales import run_extract_sales_task from tasks.sales_ingest.snowflake_check_ingested_sales import check_ingested_sales_task from tasks.sales_ingest.snowflake_load_sales import load_sales_to_staging_table_task dag = DAG( dag_id=constants.DAG_SALES_INGEST_NAME, description='Ingest sales data from StatementDB to Snowflake', start_date=datetime(2019, 1, 1), default_args={ 'provide_context': True, 'on_failure_callback': slack_failure_callback, }, schedule_interval=None ) """ ###################################################################### ############################# OPERATORS ############################## ###################################################################### """ determine_sales_ingest_type_operator = PythonOperator( dag=dag, python_callable=determine_sales_ingest_type_task, task_id='determine_sales_ingest_type' ) check_ingested_sales_operator = ShortCircuitOperator( dag=dag, python_callable=check_ingested_sales_task, task_id='check_ingested_sales' ) run_extract_sales_operator = PythonOperator( dag=dag, python_callable=run_extract_sales_task, task_id='run_extract_sales' ) load_sales_to_staging_table_operator = PythonOperator( dag=dag, python_callable=load_sales_to_staging_table_task, task_id='load_sales_to_staging_table' ) # DAG determine_sales_ingest_type_operator >> \ check_ingested_sales_operator >> \ run_extract_sales_operator >> \ load_sales_to_staging_table_operator