from datetime import datetime from airflow import DAG from airflow.models.param import Param from airflow.operators.python import PythonOperator from lib.connections import set_connections from tasks.sales_ingest.check_ingested_sales import check_ingested_sales from tasks.sales_ingest.run_extract_sales_task import run_extract_sales_task # from tasks.sales_ingest.load_sales import load_sales # from tasks.sales_ingest.copy_sales import copy_sales from tasks.sales_ingest.load_sales_from_stage import load_sales_from_stage # Set AWS and Snowflake connections set_connections() dag = DAG( 'sales_ingest', description='Ingest sales data from StatementDB to Snowflake', start_date=datetime(2019, 1, 1), default_args={'provide_context': True}, schedule_interval=None, params={ 'batch_id': Param(type='string') } ) # OPERATORS check_ingested_sales_operator = PythonOperator( dag=dag, python_callable=check_ingested_sales, task_id='check_ingested_sales' ) run_extract_sales_task_operator = PythonOperator( dag=dag, python_callable=run_extract_sales_task, task_id='run_extract_sales_task' ) # load_sales_operator = PythonOperator( # dag=dag, # python_callable=load_sales, # task_id='load_sales' # ) # # copy_sales_operator = PythonOperator( # dag=dag, # python_callable=copy_sales, # task_id='copy_sales' # ) load_sales_from_stage_operator = PythonOperator( dag=dag, python_callable=load_sales_from_stage, task_id='load_sales_from_stage' ) # DAG # check_ingested_sales_operator >> run_extract_sales_task_operator >> \ # load_sales_operator >> copy_sales_operator check_ingested_sales_operator >> run_extract_sales_task_operator >> \ load_sales_from_stage_operator