import snowflake.connector import boto3 import json import os from cryptography.hazmat.backends import default_backend from cryptography.hazmat.primitives.asymmetric import rsa from cryptography.hazmat.primitives.asymmetric import dsa from cryptography.hazmat.primitives import serialization ENVIRONMENT = os.environ.get('Environment', 'dev') OUTFILE_LOCATION = 's3://dev-royalties-sales-files/statement-db-export' SNOWFLAKE_USER = os.environ.get('SNOWFLAKE_USER') SNOWFLAKE_WAREHOUSE = os.environ.get('SNOWFLAKE_WAREHOUSE') def handler(event, context): #Handle SFN input activity_input = event.strip('\"') # Load data print('New Load') print(activity_input) region = os.environ['AWS_REGION'] sm_client = boto3.client('secretsmanager', region_name=region) sf_private_key = sm_client.get_secret_value(SecretId=ENVIRONMENT+'/statementdb-to-snowflake/SNOWFLAKE_PRIVATE_KEY')[ 'SecretString' ] sf_private_key_passphrase = sm_client.get_secret_value(SecretId=ENVIRONMENT+'/statementdb-to-snowflake/SNOWFLAKE_PRIVATE_KEY_PASSPHRASE')[ 'SecretString' ] p_key = serialization.load_pem_private_key( sf_private_key.encode(), password=sf_private_key_passphrase.encode(), backend=default_backend()) pkb = p_key.private_bytes( encoding=serialization.Encoding.DER, format=serialization.PrivateFormat.PKCS8, encryption_algorithm=serialization.NoEncryption()) ctx = snowflake.connector.connect( user=SNOWFLAKE_USER, account='orchard', private_key=pkb, database='ROYALTY_ACCOUNTING', schema='DEV', ocsp_response_cache_filename="/tmp/ocsp_response_cache" ) cs = ctx.cursor() try: sql = '''COPY INTO ROYALTY_ACCOUNTING.DEV.DIG_SALES_DETAIL FROM OUTFILE_LOCATION/OUTFILE/ CREDENTIALS=( AWS_KEY_ID='KEYKEYKEY' AWS_SECRET_KEY='SECRETSECRETSECRET' AWS_TOKEN='TOKENTOKENTOKEN') ON_ERROR=ABORT_STATEMENT FILE_FORMAT=( TYPE='CSV' FIELD_DELIMITER='\t' RECORD_DELIMITER='\n' COMPRESSION='GZIP' NULL_IF=('NULL', '0000-00-00', '2009-11-00', '2012-02-00', '2010-11-00', '2008-11-00', '2011-03-00', '0000-00-00 00:00:00', '2014-00-00') EMPTY_FIELD_AS_NULL=FALSE TRIM_SPACE=TRUE )''' sql = sql.replace("KEYKEYKEY", os.environ.get('AWS_ACCESS_KEY_ID')) sql = sql.replace("SECRETSECRETSECRET", os.environ.get('AWS_SECRET_ACCESS_KEY')) sql = sql.replace("TOKENTOKENTOKEN", os.environ.get('AWS_SESSION_TOKEN')) sql = sql.replace("OUTFILE_LOCATION", OUTFILE_LOCATION) sql = sql.replace("OUTFILE", activity_input) print("SELECTING WAREHOUSE") use_wh_sql = "USE warehouse " + SNOWFLAKE_WAREHOUSE ctx.cursor().execute(use_wh_sql) print("LOADING DATA") print(sql) cs.execute(sql) print("DATA LOAD COMPLETE") return { 'statusCode': 200, 'body': json.dumps(event) } finally: cs.close() ctx.close()