"""Helper functions related to Snowflake DB actions.""" import backoff import snowflake import sqlalchemy from integration_scripts.connectors.logging import logger from integration_scripts.connectors import snowflake as snowdb from integration_scripts.db_utils import run_query, run_bind_query, \ zip_resultproxy_dict from integration_scripts.sql.select_all_distinct_single_field \ import SELECT_ALL_DISTINCT_SINGLE_FIELD from integration_scripts.sql.select_all_single_field import \ SELECT_ALL_SINGLE_FIELD from integration_scripts.sql.create_table_as_select import \ CREATE_TABLE_AS_SELECT from integration_scripts.sql.select_all_from_table import SELECT_ALL_FROM_TABLE from integration_scripts.sql.drop_table import DROP_TABLE from integration_scripts.sql.copy_into_table_from_s3 import \ COPY_INTO_TABLE_FROM_S3 from integration_scripts.sql.drop_s3_stage import DROP_S3_STAGE @backoff.on_exception(backoff.expo, (snowflake.connector.errors.Error, sqlalchemy.exc.SQLAlchemyError), max_tries=5, max_time=1200, max_value=300) @snowdb.db_session_wrap def drop_s3_stage(session, db, schema, s3_stage_name): """Drop S3 Stage.""" db_and_schema = '{}.{}'.format(db, schema) sql = DROP_S3_STAGE params = { 'snowflake_schema': db_and_schema, 's3_stage_name': s3_stage_name } try: run_query(session, sql=sql, params=params) except Exception as e: logger.error('Drop S3 stage \'{}\' failed.'.format(s3_stage_name)) raise e @backoff.on_exception(backoff.expo, (snowflake.connector.errors.Error, sqlalchemy.exc.SQLAlchemyError), max_tries=5, max_time=1200, max_value=300) @snowdb.db_session_wrap def copy_into_table_from_s3(session, db, schema, s3_stage_name, s3_table_name): """Copy into table from S3.""" db_and_schema = '{}.{}'.format(db, schema) # Load all files from S3 stage into unified table sql = COPY_INTO_TABLE_FROM_S3 params = { 'snowflake_schema': db_and_schema, 's3_table_name': s3_table_name, 's3_stage_name': s3_stage_name } try: run_query(session, sql=sql, params=params) except Exception as e: s3_source_name = '{}.{}'.format(db_and_schema, s3_table_name) logger.error('S3 copy from \'{}\' to \'{}\' failed.'.format( s3_stage_name, s3_source_name)) raise e @snowdb.db_session_wrap def create_table_as_select(session, db, schema, source_table, target_table): """Create a table with the content from another table/view.""" db_and_schema = '{}.{}'.format(db, schema) params = { 'snowflake_schema': db_and_schema, 'source_table_name': source_table, 'target_table_name': target_table } sql = CREATE_TABLE_AS_SELECT try: run_query(session, sql=sql, params=params) except Exception as e: logger.error('Create as select \'{}\' has failed: {}'.format( target_table, str(e))) raise e @backoff.on_exception(backoff.expo, (snowflake.connector.errors.Error, sqlalchemy.exc.SQLAlchemyError), max_tries=5, max_time=1200, max_value=300) @snowdb.db_session_wrap def drop_table(session, db, schema, table_name): """Drop Source Table.""" db_and_schema = '{}.{}'.format(db, schema) sql = DROP_TABLE params = { 'snowflake_schema': db_and_schema, 'table_name': table_name } try: run_query(session, sql=sql, params=params) except Exception as e: logger.error('Drop table \'{}\' failed: {}.'.format(table_name, str(e))) raise e @snowdb.db_session_wrap def get_single_snowflake_field(session, field, table, distinct=False): """Get the value of a single field from snowflake. Args: session (SQLAlchemy): Session from db wrapper field (str): field name to retrieve table (str): table name to query distinct (bool): require distinct rows in result? Returns: dict """ params = { 'field_name': field, 'table_name': table } sql = SELECT_ALL_DISTINCT_SINGLE_FIELD if distinct \ else SELECT_ALL_SINGLE_FIELD query_results = run_query(session, sql=sql, params=params) results = zip_resultproxy_dict(query_results) return results @snowdb.db_session_wrap def get_all_rows(session, table): """Get all result rows as a dict. Args: session (SQLAlchemy): Session from db wrapper table (str): table name to query Returns: dict """ params = { 'table_name': table, } sql = SELECT_ALL_FROM_TABLE query_results = run_query(session, sql=sql, params=params) results = zip_resultproxy_dict(query_results) return results # TODO: add args, etc. to docstring @backoff.on_exception(backoff.expo, (snowflake.connector.errors.Error, sqlalchemy.exc.SQLAlchemyError), max_tries=5, max_time=1200, max_value=300) @snowdb.db_session_wrap def get_snowflake_bind_results(session, sql, params): """Return a dict fot run_query() to immediately return a dict object.""" return zip_resultproxy_dict( run_bind_query(session, sql=sql, params=params)) # TODO: add args, etc. to docstring @snowdb.db_session_wrap def get_first_value(session, sql, params=None): """Return first value for run_query() to.""" return run_query(session, sql=sql, params=params).first().items()[0][1] # TODO: add args, etc. to docstring @backoff.on_exception(backoff.expo, (snowflake.connector.errors.Error, sqlalchemy.exc.SQLAlchemyError), max_tries=5, max_time=1200, max_value=300) @snowdb.db_session_wrap def get_snowflake_results(session, sql, params=None): """Return a dict fot run_query() to immediately return a dict object.""" return zip_resultproxy_dict(run_query(session, sql=sql, params=params))