"""Ingest Weekly Spotify Market Share.""" import config import boto3 import csv import datetime import xlrd INGEST_WEEK = '24' INGEST_YEAR = '2019' CNX = config.credentials.cursor() def excel_to_csv(preprocessed_path, processed_path): """Convert xlsx to csv.""" raw_xlsx = xlrd.open_workbook(preprocessed_path) sh = raw_xlsx.sheet_by_name('Sheet1') raw_csv = open(processed_path, 'w') wr = csv.writer(raw_csv, quoting=csv.QUOTE_ALL) for rownum in range(sh.nrows): wr.writerow(sh.row_values(rownum)) raw_csv.close() def upload_csv_to_s3(processed_path, s3_path): """Upload csv from local to S3.""" s3_resource = boto3.client('s3') s3_resource.upload_file(processed_path, config.S3_BUCKET, s3_path) def create_temp_staging_raw_table(): """Create or Replace Temp Staging Table.""" sql_file = 'queries/create_temp_staging_raw_table.sql' sql = _load_query(sql_file) CNX.execute(sql.format( db=config.SNOWFLAKE_DATABASE, schema=config.SNOWFLAKE_SCHEMA)) def populate_temp_staging_raw_table(csv_filename): """Load data from s3 to temp table in Snowflake.""" sql_file = 'queries/load_temp_staging_raw.sql' sql = _load_query(sql_file) CNX.execute(sql.format( bucket=config.S3_BUCKET, dir=config.S3_DIR, filename=csv_filename, aws_key_id=config.AWS_ACCESS_KEY_ID, aws_secret_key=config.AWS_SECRET_ACCESS_KEY, db=config.SNOWFLAKE_DATABASE, schema=config.SNOWFLAKE_SCHEMA)) def populate_staging_raw_table(): """Load data from temp table to raw table.""" sql_file = 'queries/populate_staging_raw.sql' sql = _load_query(sql_file) CNX.execute(sql.format( week_start_date=_get_start_date(), ingest_week=INGEST_WEEK, ingest_year=INGEST_YEAR, db=config.SNOWFLAKE_DATABASE, schema=config.SNOWFLAKE_SCHEMA)) def _load_query(filename): """Load query.""" path = filename oq = open(path, 'r') query = oq.read() oq.close return query def _get_weekly_filename(filename): """Create complete weekly filename.""" output_filename = filename.format(week=INGEST_WEEK, year=INGEST_YEAR) return output_filename def _get_start_date(): """Find week_start_date""" week_date = '{YYYY}-W{WW}'.format(YYYY=INGEST_YEAR, WW=INGEST_WEEK) date_obj = datetime.datetime.strptime(week_date + '-1', "%Y-W%W-%w") date_obj_week_start = date_obj - datetime.timedelta(days=6) return date_obj_week_start.strftime('%Y-%m-%d') def main(): """Execute ETL.""" preprocessed_path = _get_weekly_filename(config.PREPROCESSED_PATH) csv_filename = _get_weekly_filename(config.CSV_FILENAME) processed_path = config.PROCESSED_PATH + csv_filename s3_path = '{dir}/{filename}'.format(dir=config.S3_DIR, filename=csv_filename) print('0/4 Converting xlsx to csv...') excel_to_csv(preprocessed_path, processed_path) print('1/4 Uploading file to s3...') upload_csv_to_s3(processed_path, s3_path) print('2/4 Preparing Snowflake...') create_temp_staging_raw_table() print('3/4 Loading data into Snowflake...') populate_temp_staging_raw_table(csv_filename) print('4/4 Loading data into staging_raw_spotify_weekly_ms table...') populate_staging_raw_table() if __name__ == '__main__': main()