import logging from tempfile import TemporaryDirectory import pendulum from airflow.providers.amazon.aws.hooks.s3 import S3Hook from airflow.providers.snowflake.hooks.snowflake import SnowflakeHook from flows.sme_latam import api, config logger = logging.getLogger(__name__) def bootstrap(dag_run, data_interval_start: pendulum.datetime, **kwargs): week_number = dag_run.conf.get('week', data_interval_start.week_of_year) week_year = dag_run.conf.get('year', data_interval_start.year) # TODO: discuss with provider week_year adjustments around new year logger.info(f'Week number {week_number} Year: {week_year}') format_args = dict( week_number=week_number, week_year=week_year ) drop_s3_url_template = f's3://{config.DROP_S3_BUCKET}/{config.DROP_S3_KEY_PATH_TEMPLATE}{config.DROP_FILENAME_TEMPLATE}' drop_s3_url = drop_s3_url_template.format( **format_args ) archive_s3_url_template = f's3://{config.ARCHIVE_S3_BUCKET}/{config.ARCHIVE_S3_KEY_PATH_TEMPLATE}{config.DROP_FILENAME_TEMPLATE}' archive_s3_url = archive_s3_url_template.format( **format_args ) temp_table = f'{config.TEMP_TABLE_NAME}_{week_year}W{week_number}' return { 'week_number': week_number, 'week_year': week_year, 'drop_s3_url': drop_s3_url, 'archive_s3_url': archive_s3_url, 'temp_table': temp_table, } def fetch_and_load_temp_table( drop_s3_url, temp_table_name, snowflake_conn_id='snowflake_default', **kwargs): s3_hook = S3Hook(aws_conn_id=config.AWS_CONN_ID) with TemporaryDirectory() as temp_dir: logger.info(f'Downloading {drop_s3_url} ...') downloaded_file = s3_hook.download_file( key=drop_s3_url, local_path=temp_dir, ) logger.info(f'Reading downloaded file {downloaded_file} ...') df, year, week = api.load_latam_unified_weekly_chart( xlsx_file=downloaded_file ) logger.info(f'Loaded {len(df)} rows, Year:{year}, Week:{week}') snowflake_hook = SnowflakeHook( snowflake_conn_id=snowflake_conn_id, ) logger.info(f'Saving data to {temp_table_name}') result, n_rows = api.save_dataframe_to_table( connection=snowflake_hook.get_conn(), df=df, table_name=temp_table_name ) logger.info(f'Saved {n_rows} success: {result}') return { 'rows': n_rows, 'year': year, 'week': week }