"""Lambda ppl_sales_stmt function module.""" import csv import datetime import os.path import time import urllib.parse import boto3 from botocore.exceptions import ClientError import smart_open from snowflake.ingest.simple_ingest_manager import SimpleIngestManager from snowflake.ingest import StagedFile from common_config import logger from constants import general import config 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 = urllib.parse.unquote_plus(record['s3']['object']['key']) logger.info('Bucket: ' + bucket + ' Key: ' + key) csv_rows = read_s3_file(bucket, key, record) filename = os.path.basename(key) filesize = record['s3']['object']['size'] csv_rows = update_csv_data( [row[:16] for row in csv_rows[1:]], filename, filesize) write_to_s3(bucket, key, csv_rows) 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 read_s3_file(bucket, key, record): """Read s3 CSV file and return CSV rows. Args: bucket (string): S3 bucket name key (string): S3 key name record (dict): S3 event records Returns: list: updated csv file rows """ s3 = boto3.resource('s3') obj = s3.Object(bucket, urllib.parse.unquote(key)) lines = obj.get()['Body'].read().decode('utf-8').splitlines(True) reader = csv.reader(lines) csv_rows = [row for row in reader] return csv_rows def update_csv_data(csv_rows, filename, filesize): """Update CSV with dates, filename and filesize. Args: csv_rows (list): csv file rows filename (string): CSV file name filesize (string): csv file size Returns: list: updated csv file rows """ [row.extend([ datetime.date(int(time.strptime(row[3], '%Y-%m')[0]), 1, 1).strftime( '%d-%b-%Y'), datetime.date(int(time.strptime(row[3], '%Y-%m')[0]), 12, 31).strftime( '%d-%b-%Y'), filename, filesize]) for row in csv_rows] return csv_rows def write_to_s3(bucket, key, data): """Write updated CSV with filename and filesize, into new file. Args: bucket (string): S3 bucket name key (string): S3 key name data (list): csv file rows """ filename = os.path.basename(key) source_url = 's3://{bucket_name}/{ppl_stmt_output_path}/{filename}'.format( bucket_name=bucket, ppl_stmt_output_path=general.PPL_STMT_OUTPUT_PATH, filename=filename) with smart_open.smart_open(source_url, 'w') as fileobj: writer = csv.writer(fileobj, quoting=csv.QUOTE_ALL) writer.writerows(data) 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'}