# -*- coding: utf-8 -*- """ Created on Fri May 31 15:04:32 2024 @author: CHEO001 """ import pandas as pd from sqlalchemy import create_engine from dotenv import load_dotenv import os from snowflake.sqlalchemy import URL from cryptography.hazmat.backends import default_backend from cryptography.hazmat.primitives import serialization # from cryptography.hazmat.primitives.asymmetric import rsa # from cryptography.hazmat.primitives.asymmetric import dsa load_dotenv() sfAccount = os.getenv('SFACCOUNT') sfUser = os.getenv('SFUSER') sfPswd = os.getenv('SFPSWD') sfRole = os.getenv('SFROLE') sfWarehouse = os.getenv('SFWAREHOUSE') sfPKPath = os.getenv('SFPRIVATEKEYPATH') with open(sfPKPath, "rb") as key: p_key= serialization.load_pem_private_key( key.read(), password=None , #os.environ['PRIVATE_KEY_PASSPHRASE'].encode(), backend=default_backend() ) pkb = p_key.private_bytes( encoding=serialization.Encoding.DER, format=serialization.PrivateFormat.PKCS8, encryption_algorithm=serialization.NoEncryption()) # https://docs.snowflake.com/en/developer-guide/python-connector/sqlalchemy # pip install --upgrade snowflake-sqlalchemy def getEngine(database='', schema=''): engine = create_engine(URL( account = sfAccount, user = sfUser, database = database, schema = schema, warehouse = sfWarehouse, role = sfRole ), connect_args={ 'private_key': pkb, }, ) return engine # V2 snowflake DF functions that uses SQLAlchemy def sfDF(SQLquery, sfRole=sfRole, sfWarehouse=sfWarehouse): # https://docs.snowflake.com/en/developer-guide/python-connector/sqlalchemy#opening-and-closing-a-connection try: ctx = getEngine() df = pd.read_sql(SQLquery, ctx) # converts all field names to uppercase as SQLAlchemy returns lowercases df.columns = [x.upper() for x in df.columns] except Exception as e: print(f'{e}') finally: ctx.dispose(close=True) return df # https://docs.snowflake.com/en/developer-guide/python-connector/python-connector-api#write_pandas def sfwriteDF(df, database, schema, table): try: print(f'Writing {len(df)} lines to {database}.{schema}.{table}...') database, schema, table = database.casefold(), schema.casefold(), table.casefold() ctx = getEngine(database, schema) # write_pandas(ctx, df, table) df.to_sql(table, ctx, index=False, if_exists='append') print('Completed.') except Exception as e: print(f'{e}') finally: ctx.dispose(close=True)