"""Lambda soundexchange_sales_stmt function module.""" from collections import OrderedDict import csv import datetime import os.path from urllib.parse import unquote_plus import boto3 from botocore.exceptions import ClientError from common_config import logger import config from constants import general from smart_open import smart_open from snowflake.ingest.simple_ingest_manager import SimpleIngestManager from snowflake.ingest import StagedFile from warmer_util import catch_warmer_event @catch_warmer_event() def handler(event, context): """Lambda entry point.""" try: private_key = get_snowflake_private_key() # Proxy object that abstracts the Snowpipe REST API ingest_manager = SimpleIngestManager( account=config.SNOWFLAKE_ACCOUNT, user=config.SNOWFLAKE_USER, pipe='{sf_db_name}.{sf_schema_name}.{sf_pipe_name}'.format( sf_db_name=config.SNOWFLAKE_DATABASE, sf_schema_name=config.SNOWFLAKE_SCHEMA, sf_pipe_name=config.SNOWFLAKE_PIPE_NAME), private_key=private_key) for record in event['Records']: bucket = record['s3']['bucket']['name'] key = unquote_plus(record['s3']['object']['key']) logger.info('Bucket: ' + bucket + ' Key: ' + key) filename = os.path.basename(key) filesize = record['s3']['object']['size'] write_to_s3(bucket, key, filename, filesize) logger.info('Creating files for ingest') ingest_file(event, ingest_manager, bucket, key) except Exception as e: logger.exception(str(e)) raise def get_snowflake_private_key(): """Get private key for Snowflake API call. Returns: string: Snowflake private key """ secret_name = config.AWS_SECRET_NAME endpoint_url = config.AWS_SECRET_MANAGER_URL session = boto3.session.Session() client = session.client( service_name='secretsmanager', region_name=config.AWS_REGION, endpoint_url=endpoint_url) try: get_secret_value_response = client.get_secret_value( SecretId=secret_name) except ClientError as e: if e.response['Error']['Code'] == 'ResourceNotFoundException': logger.error( 'The requested secret ' + secret_name + ' was not found') elif e.response['Error']['Code'] == 'InvalidRequestException': logger.error('The request was invalid due to:' + str(e)) elif e.response['Error']['Code'] == 'InvalidParameterException': logger.error('The request had invalid params:' + str(e)) else: logger.error( 'There is some error due to : {}'.format( str(e.response['Error']['Code']))) raise except Exception as e: logger.exception(str(e)) raise else: secret = get_secret_value_response['SecretString'] return secret def write_to_s3(bucket, key, filename, filesize): """Write updated CSV with filename and filesize, into new file. Args: bucket (string): S3 bucket name key (string): S3 key name filename (string): CSV file name filesize (string): csv file size """ filename = os.path.basename(key) input_url = 's3://{bucket}/{key}'.format(bucket=bucket, key=key) output_url = 's3://{bucket_name}/{output_path}/{filename}'.format( bucket_name=bucket, output_path=general.SOUNDEXCHANGE_STMT_OUTPUT_PATH, filename=filename) with smart_open(input_url, 'r') as infile, \ smart_open(output_url, 'w') as outfile: fieldnames = general.HEADER_ORDER writer = csv.DictWriter(outfile, fieldnames=fieldnames) writer.writeheader() new_row = {} new_rows = [] for row in csv.DictReader(l.replace('\0', '') for l in infile): for header in row: new_row[header] = row[header] date_columns = ['Broadcast Start Date', 'Broadcast End Date'] if (header in date_columns and row[header]): new_row[header] = datetime.datetime.strptime( row[header], '%d-%b-%y').strftime('%d-%b-%Y') header_mapping = general.HEADER_MAPPING for mapping_key, mapping_value in header_mapping.items(): if header == mapping_key: new_row[mapping_value] = row[mapping_key] del new_row[header] break new_row['File Name'] = filename new_row['File Size'] = filesize new_rows.append(OrderedDict(new_row)) writer.writerows(new_rows) def ingest_file(event, ingest_manager, bucket, key): """Ingest new CSV rows into Snowflake table. Args: event (dict): S3 event ingest_manager (SimpleIngestManager): ingest manager for Snowflake API bucket (string): S3 bucket name key (string): S3 key name Returns: dict: {'status': 'OK'} """ filename = os.path.basename(key) staged_file_list = [] staged_file_list.append(StagedFile(str(filename), None)) logger.info(staged_file_list) logger.info('Pushing file list to ingest REST API') response = ingest_manager.ingest_files(staged_file_list) logger.info(response) s3 = boto3.resource('s3') if response['responseCode'] == general.SUCCESS_CODE: s3.meta.client.delete_object(Bucket=bucket, Key=key) return {'status': 'OK'}