import base64 import pandas as pd import snowflake.connector import app.config as config from app.logger import log from app.secrets_manager import SecretsManagerClient from app.utils import read_query class SnowflakeClient: """Snowflake client.""" def __init__(self): log.info(f"Initializing Snowflake Client for {config.SNOWFLAKE_USER}") secrets_manager_client = SecretsManagerClient() private_key = secrets_manager_client.get_secret( config.SNOWFLAKE_PRIVATE_KEY_SECRET) try: self.conn = snowflake.connector.connect( user=config.SNOWFLAKE_USER, account=config.SNOWFLAKE_ACCOUNT, role=config.SNOWFLAKE_ROLE, private_key=base64.b64decode( bytes(private_key, encoding='utf-8')), warehouse=config.SNOWFLAKE_WAREHOUSE, ) except Exception as e: log.error("Snowflake connection failed:", e) DEBUG_LOG_MAX_LENGTH = 500 def execute_query(self, query_name, params=None, **kwargs): query_string = read_query(f'{query_name}.sql') for key, value in kwargs.items(): query_string = query_string.replace(f"${key}", value) log.info("Executing snowflake query") if len(query_string) > self.DEBUG_LOG_MAX_LENGTH: log.debug(f"{query_string[:self.DEBUG_LOG_MAX_LENGTH]}... " f"[truncated, {len(query_string)} chars total, {len(params) if params else 0} bound params]") else: log.debug(query_string) cursor = self.conn.cursor() cursor.execute(query_string, params) column_names = [desc[0] for desc in cursor.description] df = pd.DataFrame(cursor, columns=column_names) cursor.close() return df def close_connection(self): log.info("Closing snowflake connection") self.conn.close()