import os import pandas as pd import numpy as np import random import data_prep_aws import predictors import logging import time import pickle import datetime from datetime import date from imblearn.over_sampling import SMOTE from sklearn.preprocessing import OrdinalEncoder import subprocess subprocess.check_call('pip install -r /opt/ml/processing/input/dependencies/requirements-maze.rtf', shell=True) import warnings warnings.filterwarnings("ignore") import snowflake.connector 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 from password import PRIVATE_KEY_PASSPHRASE with open("/Users/impr001/Keys_snowflake/rsa_key.p8", "rb") as key: p_key= serialization.load_pem_private_key( key.read(), password=PRIVATE_KEY_PASSPHRASE.encode(), backend=default_backend() ) pkb = p_key.private_bytes( encoding=serialization.Encoding.DER, format=serialization.PrivateFormat.PKCS8, encryption_algorithm=serialization.NoEncryption()) def timing_val(func): def wrapper(*arg, **kw): t1 = time.time() res = func(*arg, **kw) t2 = time.time() return (t2 - t1), res, func.__name__ return wrapper @timing_val def load_full_data(): ctx = snowflake.connector.connect( user='eimpara', account='orchard', private_key=pkb, role= 'PROD_DATALYTICS_ROLE', warehouse = "DEV_ANALYTICS_ORCHARD" ) print("And here...") try: cs = ctx.cursor() print("Here") sql = """ select * from dev_engineering.eimpara.maze_country_table; """ cs.execute(sql) original = cs.fetch_pandas_all() return original print(original.shape) finally: cs.close() ctx.close() log_level = "DEBUG" ticket_code = "EXP_1" logging.set_verbosity(log_level) logging.debug("READY!!!") sec_id = 'dev/sagemaker-notebook-instance/SNOWFLAKE_PASSWORD' def get_secret_value(name, version=None): """Gets the value of a secret. Version (if defined) is used to retrieve a particular version of the secret. """ secrets_client = boto3.client("secretsmanager", region_name='us-east-1') kwargs = {'SecretId': name} if version is not None: kwargs['VersionStage'] = version response = secrets_client.get_secret_value(**kwargs) return response def get_snowflake_creds(username="SAGEMAKER", account="orchard", warehouse="DEV_OWS_ENGINEERING"): creds = { "user": username, "password": get_secret_value(sec_id)['SecretString'], "account": "orchard", "warehouse": warehouse, "protocol": 'https' } return creds def snowflake_connector_factory(creds=None): try: if creds: _creds = creds else: _creds = get_snowflake_creds() return snowflake.connector.connect(**_creds).cursor() except Exception as e: logging.error(f"Something went wrong - {str(e)}") def _is_version_number(s): "Check and returns true if its a version number" return re.search("^[0-9][.0-9]*[0-9]$", s) is not None def test_connection(): """ tests connection to snowflake """ with snowflake_connector_factory() as cs: try: cs.execute("SELECT current_version()") one_row = cs.fetchone() assert len(one_row) == 1 assert _is_version_number(one_row[0]) logging.info(f"Your snowflake version - {one_row[0]} PASSED!") except Exception as e: logging.error(f"Something went wrong - {str(e)}") if __name__ == '__main__': with snowflake_connector_factory(get_snowflake_creds()) as cs: try: cs.execute("USE WAREHOUSE DEV_PERFORMANCE_WAREHOUSE;") cs.execute(""" select * from dev_engineering.eimpara.maze_country_table; """) rows = cs.fetchall() except Exception as e: logging.error(f"Something went wrong - {str(e)}") data_df = pd.DataFrame(rows, columns=list(map(lambda meta: meta[0], cs.description))) df = data_df.drop_duplicates().copy() print("{} rows loaded in {} seconds".format(df.shape[0], timing)) (X_train, y_train, X_test, y_test) = data_prep_aws.split_train_test_data(df) log_model = predictors.LogisticRegressionModel(X_train, y_train) log_model.train() log_performances = predictors.estimate_predictor(log_model, X_test, y_test) logging.info(f"Perfromance metric for Logistic Model {log_performances}") matrix_log = predictors.confusion_matrix_calculation(log_model, X_test, y_test) logging.info(f"Perfromance metric for Logistic Model {matrix_log}") # with open('/Users/impr001/Documents/Jupyter_Notebooks/Maze/log_model_GB.pkl', 'wb') as file: # pickle.dump(log_model, file) # Save the model to disk pkl_filename = 'log_model_GB.pkl' with open(pkl_filename, 'wb') as file: pickle.dump(log_model, file) s3_resource = boto3.resource('s3') bucket = 'dev-cucumbers/eimpara/MAZE' key = 'log_model_GB.pkl' # Read the pickled file as bytes with open(pkl_filename, 'rb') as f: pickle_bytes_obj = f.read() # Upload the pickled file to S3 s3_resource.Bucket(bucket).put_object(Key=key, Body=pickle_bytes_obj)