"""Util functions and helpers related to databases.""" import pymysql from snowflake import connector as sf_connector from snowflake_connector.etl_connector import SnowflakeSQLExecutor from processing_accounting.conf.settings import SF_CONFIG from processing_accounting.util import setting def ResultIter(cursor, arraysize=100): """An iterator that uses fetchmany to keep memory usage down Args: cursor (postgresql cursor): """ while True: results = cursor.fetchmany(arraysize) if not results: break for result in results: yield result def _get_mysql_credentials(): """Return credentials. Returns: dict: dictionary of mysql credentials. """ return dict( host=setting.ORCHARDETL_ARTRELATIONS_HOST, port=setting.ORCHARDETL_ARTRELATIONS_PORT, user=setting.ORCHARDETL_ARTRELATIONS_USER, password=setting.ORCHARDETL_ARTRELATIONS_PASSWORD, db=setting.ORCHARDETL_ARTRELATIONS_DB) def snowflake_execute(sql, **params): """Execute a SQL query on the Snowflake DB. Args: sql (str): SQL query. ... """ with SnowflakeSQLExecutor(sf_config=SF_CONFIG) as executor: executor.execute(sql, params=params) def snowflake_query(sql): """Extract data from Snowflake. Args: sql (str): SQL statement for extracting data from Snowflake. Returns: generator: generator of dictionary """ with sf_connector.connect(**SF_CONFIG) as conn: with conn.cursor() as cursor: cursor.execute(sql) for result in ResultIter(cursor): columns = [str(c[0]) for c in cursor.description] row = { column: str(result[columns.index(column)]) for column in columns } yield row def mysql_query(sql): """Extract data from SQL Args: sql (str): SQL statement for extracting data from Snowflake Yields: dict: database records Example: { 'column1': 'paulo', 'column2': 'kuong', 'column3': 'is', 'column4': 'awesome' }, { 'column1': 'john', 'column2': 'is', 'column3': 'awesome', 'column4': 'too' } """ connection = pymysql.connect(**_get_mysql_credentials()) with connection.cursor(pymysql.cursors.DictCursor) as cursor: cursor.execute(sql) for result in ResultIter(cursor): yield result