import logging from datetime import timedelta from airflow import DAG from airflow.operators.python import PythonOperator from airflow.providers.amazon.aws.operators.s3 import S3CopyObjectOperator from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor from common.sql_tasks import SQLTemplateOperator from flows.sme_latam import config from flows.sme_latam import tasks logger = logging.getLogger(__file__) default_args = dict( owner="bi-team", execution_timeout=timedelta(minutes=30), ) with DAG( dag_id="sme_latam", default_args=default_args, schedule_interval='@weekly', start_date=config.DAG_START_DATE, catchup=True, dagrun_timeout=timedelta(minutes=60), ) as dag: bootstrap = PythonOperator( task_id='bootstrap', python_callable=tasks.bootstrap, ) sensor = S3KeySensor( task_id='sensor', bucket_key='{{ ti.xcom_pull("bootstrap").drop_s3_url }}', aws_conn_id=config.AWS_CONN_ID, mode='reschedule', poke_interval=timedelta(hours=1).total_seconds(), timeout=config.DATA_THRESHOLD.total_seconds(), ) bootstrap >> sensor fetch_and_load_temp_table = PythonOperator( task_id=f'fetch_and_load_temp_table', python_callable=tasks.fetch_and_load_temp_table, op_kwargs=dict( drop_s3_url='{{ ti.xcom_pull("bootstrap").drop_s3_url }}', temp_table_name='{{ ti.xcom_pull("bootstrap").temp_table }}' ) ) sensor >> fetch_and_load_temp_table delete_from_staging_raw_table = SQLTemplateOperator( task_id='delete_from_staging_raw_table', template='delete_from_table.sql', parameters={ 'table_name': config.STAGING_RAW_TABLE_NAME, 'where': { 'year': '{{ ti.xcom_pull("bootstrap").week_year }}', 'week': '{{ ti.xcom_pull("bootstrap").week_number }}', } } ) fetch_and_load_temp_table >> delete_from_staging_raw_table load_staging_raw_table = SQLTemplateOperator( task_id='load_staging_raw_table', template='insert_from_select.sql', parameters={ 'source_table_name': '{{ ti.xcom_pull("bootstrap").temp_table }}', 'target_table_name': config.STAGING_RAW_TABLE_NAME, } ) delete_from_staging_raw_table >> load_staging_raw_table archive = S3CopyObjectOperator( task_id='archive', source_bucket_key='{{ ti.xcom_pull("bootstrap").drop_s3_url }}', dest_bucket_key='{{ ti.xcom_pull("bootstrap").archive_s3_url }}', aws_conn_id=config.AWS_CONN_ID ) load_staging_raw_table >> archive drop_temp_table = SQLTemplateOperator( task_id='drop_temp_table', template='drop_table.sql', parameters=dict( table_name='{{ ti.xcom_pull("bootstrap").temp_table }}' ) ) archive >> drop_temp_table