"""Helper functions related to DB actions.""" import json import sys import traceback from contextlib import contextmanager import backoff import pymysql import snowflake import sqlalchemy from integration_scripts import logger as log from integration_scripts.connectors import mysql from integration_scripts.connectors import snowflake as snowdb from constants import globals from constants import mysql as mysql_const from constants import database as db_const from constants import s3 as s3_const from constants.ingestion_key_id import foreign_key_field from integration_scripts import db_utils as orch_db from integration_scripts.connectors.snowflake import snow_db_config from integration_scripts.connectors.mysql import rds_mysl_db_config from util.exceptions import NewArtistError from util.logic_utils import slice_off_request_code from sql.snowflake.create_release_response_log import \ CREATE_RELEASE_RESPONSE_LOG as CREATE_RELEASE_RESPONSE_LOG_SNOW from sql.mysql.create_release_response_log import \ CREATE_RELEASE_RESPONSE_LOG as CREATE_RELEASE_RESPONSE_LOG_MYSQL from sql.snowflake.create_track_response_log import \ CREATE_TRACK_RESPONSE_LOG as CREATE_TRACK_RESPONSE_LOG_SNOW from sql.mysql.create_track_response_log import \ CREATE_TRACK_RESPONSE_LOG as CREATE_TRACK_RESPONSE_LOG_MYSQL from sql.snowflake.insert_release_response_log import \ INSERT_RELEASE_RESPONSE_LOG as INSERT_RELEASE_RESPONSE_LOG_SNOW from sql.snowflake.select_count_of_rows import SELECT_COUNT_OF_ROWS from sql.mysql.insert_release_response_log import \ INSERT_RELEASE_RESPONSE_LOG as INSERT_RELEASE_RESPONSE_LOG_MYSQL from sql.snowflake.insert_track_response_log import \ INSERT_TRACK_RESPONSE_LOG as INSERT_TRACK_RESPONSE_LOG_SNOW from sql.mysql.insert_track_response_log import \ INSERT_TRACK_RESPONSE_LOG as INSERT_TRACK_RESPONSE_LOG_MYSQL from sql.mysql.select_matching_artists_in_label import \ SELECT_MATCHING_ARTISTS_IN_LABEL from sql.snowflake.copy_into_s3_from_join import COPY_INTO_S3_FROM_JOIN from sql.snowflake.copy_into_s3_from_join_by_label import \ COPY_INTO_S3_FROM_JOIN_BY_LABEL from sql.snowflake.create_s3_source_table import CREATE_S3_SOURCE_TABLE from sql.reporting.error_breakdown import ERROR_BREAKDOWN from sql.reporting.ingestion_action_breakdown import INGESTION_ACTION_BREAKDOWN from sql.reporting.release_level_summary import RELEASE_LEVEL_SUMMARY from sql.reporting.error_breakdown_by_label import ERROR_BREAKDOWN_BY_LABEL from sql.reporting.ingestion_action_breakdown_by_label \ import INGESTION_ACTION_BREAKDOWN_BY_LABEL from sql.reporting.release_level_summary_by_label \ import RELEASE_LEVEL_SUMMARY_BY_LABEL from sql.templates.create_bulk_upload_template_table import \ CREATE_BULK_UPLOAD_TEMPLATE_TABLE from sql.snowflake.create_audio_asset_ingestion_table import \ CREATE_AUDIO_ASSET_INGESTION_TABLE from sql.snowflake.create_cover_asset_ingestion_table import \ CREATE_COVER_ASSET_INGESTION_TABLE from sql.snowflake.create_audio_asset_ingestion_file import \ CREATE_AUDIO_ASSET_INGESTION_FILE from sql.snowflake.create_cover_asset_ingestion_file import \ CREATE_COVER_ASSET_INGESTION_FILE from sql.snowflake.create_s3_csv_input_file_format import \ CREATE_S3_CSV_INPUT_FILE_FORMAT from sql.snowflake.create_s3_csv_output_file_format import \ CREATE_S3_CSV_OUTPUT_FILE_FORMAT from sql.snowflake.create_s3_stage import \ CREATE_S3_STAGE from sql.snowflake.select_release_id_from_log \ import SELECT_RELEASE_ID_FROM_LOG from sql.snowflake.select_project_name_from_sme_labelcopy_by_upc import \ SELECT_PROJECT_NAME_FROM_SME_LABELCOPY_BY_UPC from sql.snowflake.select_product_code_from_session_id import \ SELECT_PRODUCT_CODE_FROM_SESSION_ID import config class NullFormatter(dict): """Replace missing dict keys with NULL.""" def __missing__(self, key): """Handle missing key.""" return None @contextmanager def global_semaphore(): globals.g_lock and globals.g_lock.acquire() try: yield finally: globals.g_lock and globals.g_lock.release() # Custom callback function to log a message def log_backoff(details): """Log a message when backoff is triggered. Args: details (dict): Details of the backoff event """ exception = details['exception'] exception_message = getattr(exception, 'message', str(exception)) func_name = details['target'].__name__ + '()' log.warning(f"{exception_message}") log.warning(f"{details['tries']}: {func_name} will retry after " f"{details['wait']:.2f} seconds.") @backoff.on_exception(backoff.expo, (pymysql.err.Error, sqlalchemy.exc.SQLAlchemyError, NewArtistError), base=5, factor=2, max_tries=5, max_time=1200, max_value=300) @mysql.rds_db_session_wrap def insert_mysql_release_log(session, log_table, session_id='', request_num=None, action_type='', orchard_vendor_id=None, release_id=None, product_code='', project_id=None, project_code='', request='', response=''): """Insert a release log entry. Args: session (SQLAlchemy): Session from db wrapper log_table (str): Name of log_table to insert into session_id (str): uuid hex generated by script at runtime request_num (int): The request chunk number action_type (str): The type of operation orchard_vendor_id (int): Orchard Vendor ID release_id (int): Orchard release_id product_code (str): The product code project_id (int): The project id project_code (str): The project code request (str): A str representation of the JSON request object from an ingestion attempt response (str): A str representation of the JSON response object from an ingestion attempt Returns: dict """ if not config.LOG_REQUEST: request = '' sql = INSERT_RELEASE_RESPONSE_LOG_MYSQL log.bind(multi=True).info('Inserting release log entry.') sql = sql.format( database=rds_mysl_db_config['database'], log_table=log_table.upper()) params = { 'session_id': session_id, 'request_num': request_num, 'action_type': action_type, 'orchard_vendor_id': orchard_vendor_id, 'release_id': release_id, 'product_code': product_code, 'project_id': project_id, 'project_code': project_code, 'request': json.dumps(request, ensure_ascii=False), # Can be blank 'response': json.dumps(response, ensure_ascii=False) } # globals.g_lock and globals.g_lock.acquire() # Wrap with sempahore lock try: with global_semaphore(): query_results = orch_db.run_bind_query( session, sql=sql, params=params) except Exception as e: log.bind(multi=True).error( '\n- - - - - - - - - - - - - - - - - - - - - - -\n' 'Error Writing to MySQL DB:\n' 'Single Release.\n' f'{str(e.__class__)}: {str(e)}\n' 'Traceback:\n' f'{traceback.format_exc()}\n' '- - - - - - - - - - - - - - - - - - - - - - -') if config.DEBUG_WRITE_TO_FILE: log.bind(json_only=True, multi=True).info(params) params = { 'session_id': session_id, 'request_num': request_num, 'action_type': 'log error', 'orchard_vendor_id': orchard_vendor_id, 'release_id': None, 'product_code': None, 'project_id': None, 'project_code': None, 'request': None, 'response': str(e) } with global_semaphore(): query_results = orch_db.run_bind_query( session, sql=sql, params=params) log.bind(multi=True).info('Release log error logged.') # globals.g_lock and globals.g_lock.release() results = query_results return results @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 insert_snowflake_release_log(session, log_table, session_id='', request_num=None, action_type='', orchard_vendor_id=None, release_id=None, product_code='', project_id=None, project_code='', request='', response=''): """Insert a release log entry. Args: session (SQLAlchemy): Session from db wrapper log_table(str): Name of log_table to insert into session_id (str): uuid hex generated by script at runtime request_num (int): The request chunk number action_type (str): The type of operation orchard_vendor_id (int): Orchard Vendor ID release_id (int): Orchard release_id product_code (str): The product code project_id (int): The project id project_code (str): The project code request (str): A str representation of the JSON request object from an ingestion attempt response (str): A str representation of the JSON response object from an ingestion attempt Returns: dict """ if not config.LOG_REQUEST: request = '' sql = INSERT_RELEASE_RESPONSE_LOG_SNOW sql = sql.format( snowflake_schema=snow_db_config['schema'], snowflake_database=snow_db_config['database'], log_table=log_table) params = { 'session_id': session_id, 'request_num': request_num, 'action_type': action_type, 'orchard_vendor_id': orchard_vendor_id, 'release_id': release_id, 'product_code': product_code, 'project_id': project_id, 'project_code': project_code, 'request': json.dumps(request, ensure_ascii=False), # Can be blank 'response': json.dumps(response, ensure_ascii=False) } try: with global_semaphore(): query_results = orch_db.run_bind_query( session, sql=sql, params=params) except Exception as e: log.bind(multi=True).error( '\n- - - - - - - - - - - - - - - - - - - - - - -\n' 'Error Writing to Snowflake DB:\n' 'Single Release.\n' f'{str(e.__class__)}: {str(e)}' 'Traceback:\n' f'{traceback.format_exc()}\n' '- - - - - - - - - - - - - - - - - - - - - - -') if config.DEBUG_WRITE_TO_FILE: log.bind(json_only=True, multi=True).info(params) params = { 'session_id': session_id, 'request_num': request_num, 'action_type': 'log error', 'orchard_vendor_id': orchard_vendor_id, 'release_id': None, 'product_code': None, 'project_id': None, 'project_code': None, 'request': None, 'response': str(e) } with global_semaphore(): query_results = orch_db.run_bind_query( session, sql=sql, params=params) log.bind(multi=True).info('Release log error logged.') results = query_results return results @mysql.rds_db_session_wrap def insert_multi_mysql_release_logs(session, release_log_params): """Insert a list of release log entries. Args: session (SQLAlchemy): Session from db wrapper release_log_params (list): List of insert param dicts Returns: dict """ # Final result list results = list() log_table = db_const.RELEASE_LOG_TABLE.upper() for release_log in release_log_params: # slice off request_code release_log = slice_off_request_code(release_log) # Fill Missing values with NULL release_log = NullFormatter(release_log) query_results = insert_mysql_release_log( session=session, log_table=log_table, **release_log) results.append(query_results) return results @snowdb.db_session_wrap def insert_multi_snowflake_release_logs(session, release_log_params): """Insert a list of release log entries. Args: session (SQLAlchemy): Session from db wrapper release_log_params (list): List of insert param dicts Returns: dict """ # Final result list results = list() log_table = db_const.RELEASE_LOG_TABLE.upper() for release_log in release_log_params: # slice off request_code release_log = slice_off_request_code(release_log) # Fill Missing values with NULL release_log = NullFormatter(release_log) query_results = insert_snowflake_release_log( session=session, log_table=log_table, **release_log) results.append(query_results) return results @backoff.on_exception(backoff.expo, (pymysql.err.Error, sqlalchemy.exc.SQLAlchemyError, NewArtistError), base=5, factor=2, max_tries=5, max_time=1200, max_value=300) @mysql.rds_db_session_wrap def insert_mysql_track_log(session, session_id='', request_num=None, orchard_vendor_id=None, product_code='', project_code='', track_num=None, vol_num=None, errors=''): """Insert a track log entry. Args: session: (SQLAlchemy) Session from db wrapper session_id (str): UUID hex generated by script at runtime request_num (int): The request chunk number orchard_vendor_id (int): Orchard Vendor ID product_code (str): The product code project_code (str): The project code track_num (int): Track position on volume vol_num (int): Volume number in product errors (str): List of errors Returns: dict """ sql = INSERT_TRACK_RESPONSE_LOG_MYSQL sql = sql.format( database=rds_mysl_db_config['database'], log_table=db_const.TRACK_LOG_TABLE.upper()) params = { 'session_id': session_id, 'request_num': request_num, 'orchard_vendor_id': orchard_vendor_id, 'product_code': product_code, 'project_code': project_code, 'track_num': track_num, 'vol_num': vol_num, 'errors': json.dumps(errors, ensure_ascii=False) } try: with global_semaphore(): query_results = orch_db.run_bind_query( session, sql=sql, params=params) except Exception as e: log.bind(multi=True).error( '\n- - - - - - - - - - - - - - - - - - - - - - -\n' 'Error Writing to DB:\n' 'Single Track\n' f'{str(e.__class__)}: {str(e)}' 'Traceback:\n' f'{traceback.format_exc()}\n' '- - - - - - - - - - - - - - - - - - - - - - -') if config.DEBUG_WRITE_TO_FILE: log.bind(json_only=True, multi=True).info(params) params = { 'session_id': session_id, 'request_num': request_num, 'orchard_vendor_id': orchard_vendor_id, 'product_code': product_code, 'project_code': project_code, 'track_num': track_num, 'vol_num': vol_num, 'errors': 'log error: Error Writing to DB: Single Track' } with global_semaphore(): query_results = orch_db.run_bind_query( session, sql=sql, params=params) log.bind(multi=True).info('Track log error logged.') results = query_results return results @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 insert_snowflake_track_log(session, session_id='', request_num=None, orchard_vendor_id=None, product_code='', project_code='', track_num=None, vol_num=None, errors=''): """Insert a track log entry. Args: session: (SQLAlchemy) Session from db wrapper session_id (str): UUID hex generated by script at runtime request_num (int): The request chunk number orchard_vendor_id (int): Orchard Vendor ID product_code (str): The product code project_code (str): The project code track_num (int): Track position on volume vol_num (int): Volume number in product errors (str): List of errors Returns: dict """ sql = INSERT_TRACK_RESPONSE_LOG_SNOW # TODO: RDS Feature Flag log_table = '{}_{}'.format( db_const.TRACK_LOG_TABLE.upper(), config.ENVIRONMENT.upper()) sql = sql.format( snowflake_schema=snow_db_config['schema'], snowflake_database=snow_db_config['database'], log_table=log_table) params = { 'session_id': session_id, 'request_num': request_num, 'orchard_vendor_id': orchard_vendor_id, 'product_code': product_code, 'project_code': project_code, 'track_num': track_num, 'vol_num': vol_num, 'errors': json.dumps(errors, ensure_ascii=False) } try: with global_semaphore(): query_results = orch_db.run_bind_query( session, sql=sql, params=params) except Exception as e: log.bind(multi=True).error( '\n- - - - - - - - - - - - - - - - - - - - - - -\n' 'Error Writing to DB:\n' 'Single Track\n' f'{str(e.__class__)}: {str(e)}' 'Traceback:\n' f'{traceback.format_exc()}\n' '- - - - - - - - - - - - - - - - - - - - - - -') if config.DEBUG_WRITE_TO_FILE: log.bind(json_only=True, multi=True).info(params) params = { 'session_id': session_id, 'request_num': request_num, 'orchard_vendor_id': orchard_vendor_id, 'product_code': product_code, 'project_code': project_code, 'track_num': track_num, 'vol_num': vol_num, 'errors': 'log error: Error Writing to DB: Single Track' } with global_semaphore(): query_results = orch_db.run_bind_query( session, sql=sql, params=params) log.bind(multi=True).info('Track log error logged.') results = query_results return results @mysql.rds_db_session_wrap def insert_multi_mysql_track_logs(session, track_log_params): """Insert a list of release log entries. Args: session (SQLAlchemy): Session from db wrapper track_log_params (list): List of insert param dicts Returns: dict """ # Final result list results = list() for track_log in track_log_params: # Fill Missing values with NULL track_log = NullFormatter(track_log) query_results = insert_mysql_track_log( session=session, **track_log) results.append(query_results) return results @snowdb.db_session_wrap def insert_multi_snowflake_track_logs(session, track_log_params): """Insert a list of release log entries. Args: session (SQLAlchemy): Session from db wrapper track_log_params (list): List of insert param dicts Returns: dict """ # Final result list results = list() for track_log in track_log_params: # Fill Missing values with NULL track_log = NullFormatter(track_log) query_results = insert_snowflake_track_log( session=session, **track_log) results.append(query_results) return results @snowdb.db_session_wrap def write_to_snowflake_log(session, **kwargs): """Write all the rows to log. Args: session: (SQLAlchemy) Session from db wrapper A dict of lists whose keys represent ingestion statuses in the database. """ # TODO: Figure out how to get full track list. for status, release_list in kwargs.items(): if status == 'tracks': track_log_list = release_list if len(track_log_list): log.bind(multi=True).info('Logging tracks...') # # Insert all track logs for track_params in track_log_list: insert_multi_snowflake_track_logs( session=session, track_log_params=track_params) # track count track_log_count = len(track_params) log.bind(multi=True).info( '{} track response(s) logged'.format( track_log_count)) else: insert_multi_snowflake_release_logs( session=session, release_log_params=release_list) log.bind(multi=True).info('{} releases logged as {}.'.format( len(release_list), status.replace('_', ' '))) @mysql.rds_db_session_wrap def write_to_mysql_log(session, **kwargs): """Write all the rows to log. Args: session: (SQLAlchemy) Session from db wrapper A dict of lists whose keys represent ingestion statuses in the database. """ # TODO: Figure out how to get full track list. for status, release_list in kwargs.items(): if status == 'tracks': track_log_list = release_list if len(track_log_list): log.bind(multi=True).info('Logging tracks...') # # Insert all track logs for track_params in track_log_list: insert_multi_mysql_track_logs( session=session, track_log_params=track_params) # track count track_log_count = len(track_params) log.bind(multi=True).info( '{} track response(s) logged'.format( track_log_count)) else: insert_multi_mysql_release_logs( session=session, release_log_params=release_list) log.bind(multi=True).info('{} releases logged as {}.'.format( len(release_list), status.replace('_', ' '))) @backoff.on_exception(backoff.expo, (pymysql.err.Error, sqlalchemy.exc.SQLAlchemyError, NewArtistError), base=5, factor=2, max_tries=5, max_time=1200, max_value=300) @mysql.rds_db_session_wrap def create_rds_release_log_table(session): """Create the log table if needed.""" log.info('Creating RDS release log table.') sql = CREATE_RELEASE_RESPONSE_LOG_MYSQL params = { 'database': rds_mysl_db_config['database'] } return orch_db.run_query(session, sql=sql, params=params) @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 create_snowflake_release_log_table(session): """Create the log table if needed.""" log.info('Creating Snowflake release log table.') sql = CREATE_RELEASE_RESPONSE_LOG_SNOW params = { 'snowflake_database': snow_db_config['database'], 'snowflake_schema': snow_db_config['schema'], 'environment': config.ENVIRONMENT.lower() } return orch_db.run_query(session, sql=sql, params=params) @backoff.on_exception(backoff.expo, (pymysql.err.Error, sqlalchemy.exc.SQLAlchemyError, NewArtistError), base=5, factor=2, max_tries=5, max_time=1200, max_value=300) @mysql.rds_db_session_wrap def create_rds_track_log_table(session): """Create the log table if needed.""" log.info('Creating RDS track log table.') sql = CREATE_TRACK_RESPONSE_LOG_MYSQL params = { 'database': rds_mysl_db_config['database'] } return orch_db.run_query(session, sql=sql, params=params) @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 create_snowflake_track_log_table(session): """Create the log table if needed.""" log.info('Creating Snowflake track log table.') sql = CREATE_TRACK_RESPONSE_LOG_SNOW params = { 'snowflake_database': snow_db_config['database'], 'snowflake_schema': snow_db_config['schema'], 'environment': config.ENVIRONMENT.lower() } return orch_db.run_query(session, sql=sql, params=params) @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): """Return a list of dicts from run_query() result.""" return orch_db.zip_resultproxy_dict( orch_db.run_query(session, sql=sql, params=params)) @backoff.on_exception(backoff.expo, (pymysql.err.Error, sqlalchemy.exc.SQLAlchemyError, NewArtistError), base=5, factor=2, max_tries=5, max_time=120, max_value=300, on_backoff=log_backoff) @mysql.ar_db_session_wrap def get_new_input_artists(session, label_id, input_artist_list, except_if_found=False): field_name = mysql_const.RELEASE_ARTIST_NAME_FIELD params = { 'label_id': label_id, 'artist_name_list': list(input_artist_list) } # Find all artists in the input file that are in `artist_info` already sql = SELECT_MATCHING_ARTISTS_IN_LABEL query_results = orch_db.zip_resultproxy_dict( orch_db.run_bind_query(session, sql=sql, params=params)) # Build a SET of these artists matched_artist_info_list = { m[field_name].strip().lower() for m in query_results } # Filter empty vals matched_artist_info_list = { m for m in matched_artist_info_list if m } # Get the SET of artists that are not in the artist_info table matched_artist_info_list = \ set(input_artist_list) - matched_artist_info_list if except_if_found and len(matched_artist_info_list): raise NewArtistError( message='Artists have not all been written.', result=matched_artist_info_list) return matched_artist_info_list @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 create_s3_source_table(session, s3_table_name): """Create S3 Source Table.""" sql = CREATE_S3_SOURCE_TABLE params = { 'snowflake_database': snow_db_config['database'], 'snowflake_schema': snow_db_config['schema'], 's3_table_name': s3_table_name } try: orch_db.run_query(session, sql=sql, params=params) except Exception as e: db_and_schema = snow_db_config['database'] + '.' + snow_db_config[ 'schema'] s3_source_name = db_and_schema + '.' + s3_table_name log.error( 'S3 source table \'{}\' creation failed.'.format( s3_source_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 create_upload_template(session): """Create a blank table with bulk upload schema""" sql = CREATE_BULK_UPLOAD_TEMPLATE_TABLE.format( snowflake_database=snow_db_config['database'], snowflake_schema=snow_db_config['schema'] ) try: orch_db.run_query(session, sql=sql) except Exception as e: log.error('Create Bulk Upload Template failed.') 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 join_log_and_unload_label_data( session, session_id, s3_stage_name, s3_table_name, input_file_root, label_id, track_log_table, code): """Output the result of one label joined with the run log to S3.""" # Get datetime to create unique stage and source for operation env = config.ENVIRONMENT sig = input_file_root + '_' + code ing_out_file_name = \ '{}/{}/{}/{}/{}_ingestion_output_{}_{}_{}.csv'.format( env, session_id, sig, label_id, input_file_root, code, label_id, env) sql = COPY_INTO_S3_FROM_JOIN_BY_LABEL params = { 'source_snowflake_schema': snow_db_config['schema'], 'source_snowflake_database': snow_db_config['database'], 'log_snowflake_schema': config.SNOWFLAKE_SYNC_TARGET_SCHEMA, 'log_snowflake_database': config.SNOWFLAKE_SYNC_TARGET_DB, 's3_stage_name': s3_stage_name, 'track_log_table': track_log_table, 'session_id': session_id, 'source_table': s3_table_name, 'project_code_field': db_const.project_code, 'product_code_field': db_const.product_code, 'file_name': ing_out_file_name, 'label_id': label_id } try: orch_db.run_query(session, sql=sql, params=params) except Exception as e: log.error('COPY INTO 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 join_log_and_unload_all_data( session, session_id, s3_stage_name, s3_table_name, input_file_root, track_log_table, code): """Output the result of the whole dataset joined with the run log to S3.""" # Get datetime to create unique stage and source for operation env = config.ENVIRONMENT sig = input_file_root + '_' + code ing_out_file_name = '{}/{}/{}/{}_ingestion_output_{}_{}.csv'.format( env, session_id, sig, input_file_root, code, env) sql = COPY_INTO_S3_FROM_JOIN params = { 'source_snowflake_schema': snow_db_config['schema'], 'source_snowflake_database': snow_db_config['database'], 'log_snowflake_schema': config.SNOWFLAKE_SYNC_TARGET_SCHEMA, 'log_snowflake_database': config.SNOWFLAKE_SYNC_TARGET_DB, 's3_stage_name': s3_stage_name, 'track_log_table': track_log_table, 'session_id': session_id, 'source_table': s3_table_name, 'project_code_field': db_const.project_code, 'product_code_field': db_const.product_code, 'file_name': ing_out_file_name, } try: orch_db.run_query(session, sql=sql, params=params) except Exception as e: log.error( 'COPY INTO 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 join_log_and_unload_label_reports( session, session_id, s3_stage_name, input_file_root, label_id, release_log_table, track_log_table, code): """Output the label-level summary reports to S3.""" schema_params = { 'source_snowflake_schema': snow_db_config['schema'], 'source_snowflake_database': snow_db_config['database'], 'log_snowflake_schema': config.SNOWFLAKE_SYNC_TARGET_SCHEMA, 'log_snowflake_database': config.SNOWFLAKE_SYNC_TARGET_DB, } # Get datetime to create unique stage and source for operation env = config.ENVIRONMENT sig = input_file_root + '_' + code rls_sum_file_name = \ '{}/{}/{}/{}/{}_release_level_summary_{}_{}_{}.csv'.format( env, session_id, sig, label_id, input_file_root, code, label_id, env) sql = RELEASE_LEVEL_SUMMARY_BY_LABEL params = { **schema_params, 's3_stage_name': s3_stage_name, 'file_name': rls_sum_file_name, 'release_log_table': release_log_table, 'session_id': session_id, 'label_id': label_id } try: orch_db.run_query(session, sql=sql, params=params) except Exception as e: log.error( 'Release Level Summary: COPY INTO S3 stage \'{}\' ' 'failed - {}.'.format(s3_stage_name, str(e))) else: log.info('Release Level Summary Written: {}.'.format( rls_sum_file_name)) sql = INGESTION_ACTION_BREAKDOWN_BY_LABEL ing_act_file_name = \ '{}/{}/{}/{}/{}_ingestion_action_breakdown_{}_{}_{}.csv'.format( env, session_id, sig, label_id, input_file_root, code, label_id, env) params = { **schema_params, 's3_stage_name': s3_stage_name, 'file_name': ing_act_file_name, 'release_log_table': release_log_table, 'session_id': session_id, 'label_id': label_id } try: orch_db.run_query(session, sql=sql, params=params) except Exception as e: log.error( 'Ingestion Action Breakdown: COPY INTO S3 ' 'stage \'{}\' failed - {}.'.format(s3_stage_name, str(e))) else: log.info( 'Ingestion Action Breakdown Written: {}.'.format( ing_act_file_name)) sql = ERROR_BREAKDOWN_BY_LABEL err_brk_file_name = '{}/{}/{}/{}/{}_error_breakdown_{}_{}_{}.csv'.format( env, session_id, sig, label_id, input_file_root, code, label_id, env) params = { **schema_params, 's3_stage_name': s3_stage_name, 'file_name': err_brk_file_name, 'track_log_table': track_log_table, 'session_id': session_id, 'label_id': label_id } try: orch_db.run_query(session, sql=sql, params=params) except Exception as e: log.error( 'Error Breakdown: COPY INTO S3 stage \'{}\' ' 'failed - {}.'.format( s3_stage_name, str(e))) else: log.info( 'Error Breakdown Written: {}.'.format( err_brk_file_name)) @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 join_log_and_unload_all_reports( session, session_id, s3_stage_name, input_file_root, source_table_name, release_log_table, track_log_table, code, start_time): """Output the label-level summary reports to S3.""" # Get datetime to create unique stage and source for operation time_now = start_time env = config.ENVIRONMENT schema_params = { 'source_snowflake_schema': snow_db_config['schema'], 'source_snowflake_database': snow_db_config['database'], 'log_snowflake_schema': config.SNOWFLAKE_SYNC_TARGET_SCHEMA, 'log_snowflake_database': config.SNOWFLAKE_SYNC_TARGET_DB, } sig = input_file_root + '_' + code rls_sum_file_name = \ '{}/{}/{}/{}_release_level_summary_{}_{}.csv'.format( env, session_id, sig, input_file_root, code, env) sql = RELEASE_LEVEL_SUMMARY params = { **schema_params, 's3_stage_name': s3_stage_name, 'file_name': rls_sum_file_name, 'release_log_table': release_log_table, 'session_id': session_id, } try: orch_db.run_query(session, sql=sql, params=params) except Exception as e: log.error( 'Release Level Summary: COPY INTO S3 ' 'stage \'{}\' failed - {}.'.format(s3_stage_name, str(e))) else: log.info( 'Release Level Summary Written: {}.'.format(rls_sum_file_name)) # TODO - Add / Change SQL to allow ability to process one label # across multiple documents. sql = INGESTION_ACTION_BREAKDOWN ing_act_file_name = \ '{}/{}/{}/{}_ingestion_action_breakdown_{}_{}.csv'.format( env, session_id, sig, input_file_root, code, env) params = { **schema_params, 's3_stage_name': s3_stage_name, 'file_name': ing_act_file_name, 'release_log_table': release_log_table, 'session_id': session_id, } try: orch_db.run_query(session, sql=sql, params=params) except Exception as e: log.error( 'Ingestion Action Breakdown: COPY INTO S3 ' 'stage \'{}\' failed - {}.'.format(s3_stage_name, str(e))) else: log.info( 'Ingestion Action Breakdown Written: {}.'.format( ing_act_file_name)) sql = ERROR_BREAKDOWN err_brk_file_name = '{}/{}/{}/{}_error_breakdown_{}_{}.csv'.format( env, session_id, sig, input_file_root, code, env) params = { **schema_params, 's3_stage_name': s3_stage_name, 'file_name': err_brk_file_name, 'track_log_table': track_log_table, 'session_id': session_id, } try: orch_db.run_query(session, sql=sql, params=params) except Exception as e: log.error( 'Error Breakdown: COPY INTO S3 stage \'{}\' ' 'failed - {}.'.format( s3_stage_name, str(e))) else: log.info( 'Error Breakdown Written: {}.'.format( err_brk_file_name)) sql = CREATE_AUDIO_ASSET_INGESTION_TABLE # Get datetime to create unique stage and source for operation audio_asset_table_name = \ s3_const.audio_asset_output_template + '{}'.format(time_now) params = { **schema_params, 'target_table_name': audio_asset_table_name, 'foreign_key_field': foreign_key_field, 'source_table_name': source_table_name, 'session_id': session_id } try: orch_db.run_query(session, sql=sql, params=params) except Exception as e: log.error( 'Audio Asset Ingestion Table: CREATE OR REPLACE TABLE \'{}\' ' 'failed - {}.'.format( audio_asset_table_name, str(e))) else: log.info( 'Audio Asset Ingestion Table Created: {}.'.format( audio_asset_table_name)) sql = CREATE_AUDIO_ASSET_INGESTION_FILE audio_asset_file_name = \ '{}/{}/{}/{}_audio_asset_ingestion_{}_{}.csv'.format( env, session_id, sig, input_file_root, code, env) params = { **schema_params, 's3_stage_name': s3_stage_name, 'file_name': audio_asset_file_name, 'foreign_key_field': foreign_key_field, 'source_table_name': source_table_name, 'session_id': session_id } try: orch_db.run_query(session, sql=sql, params=params) except Exception as e: log.error( 'Audio Asset Ingestion List file \'{}\' ' 'failed - {}.'.format( audio_asset_file_name, str(e))) else: log.info( 'Audio Asset Ingestion List Written: {}.'.format( audio_asset_file_name)) sql = CREATE_COVER_ASSET_INGESTION_TABLE # Get datetime to create unique stage and source for operation cover_asset_table_name = \ s3_const.cover_asset_output_template + '{}'.format(time_now) params = { **schema_params, 'target_table_name': cover_asset_table_name, 'foreign_key_field': foreign_key_field, 'source_table_name': source_table_name, 'session_id': session_id } try: orch_db.run_query(session, sql=sql, params=params) except Exception as e: log.error( 'Cover Asset Ingestion Table: CREATE OR REPLACE TABLE \'{}\' ' 'failed - {}.'.format( cover_asset_table_name, str(e))) else: log.info( 'Cover Asset Ingestion Table Created: {}.'.format( cover_asset_table_name)) sql = CREATE_COVER_ASSET_INGESTION_FILE cover_asset_file_name = \ '{}/{}/{}/{}_cover_asset_ingestion_{}_{}.csv'.format( env, session_id, sig, input_file_root, code, env) params = { **schema_params, 's3_stage_name': s3_stage_name, 'file_name': cover_asset_file_name, 'foreign_key_field': foreign_key_field, 'source_table_name': source_table_name, 'session_id': session_id } try: orch_db.run_query(session, sql=sql, params=params) except Exception as e: log.error( 'Cover Asset Ingestion List file \'{}\' ' 'failed - {}.'.format( audio_asset_file_name, str(e))) else: log.info( 'Cover Asset Ingestion List Written: {}.'.format( cover_asset_file_name)) return audio_asset_table_name, cover_asset_table_name @snowdb.db_session_wrap def select_count_of_snowflake_rows_by_session( session, db, schema, table_name, session_id): """Count rows a given Snowflake table with a session_id.""" params = { 'snowflake_schema': schema, 'snowflake_database': db, 'table_name': table_name, 'session_id': session_id } sql = SELECT_COUNT_OF_ROWS try: query_results = orch_db.run_query(session, sql=sql, params=params) except Exception as e: log.error( 'Could not count rows in {}.{}.{}: {}'.format( db, schema, table_name, str(e))) return 0 results = orch_db.zip_resultproxy_dict(query_results) if len(results): return results[0]['row_count'] return 0 # TODO: Convert all db, schema args to kwargs with snowdb.snow_db_config @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 create_s3_stage(session, bucket, path, db, schema, s3_stage_name, file_format, aws_key, aws_secret): """Create S3 Stage.""" # Join the bucket and the path to files bucket_and_path = bucket + '/' + path # Pass params and run query sql = CREATE_S3_STAGE params = { 'snowflake_schema': schema, 'snowflake_database': db, 's3_stage_name': s3_stage_name, 'file_format': file_format, 'bucket_and_path': bucket_and_path, 'aws_key_id': aws_key, 'aws_secret_key': aws_secret } try: orch_db.run_query(session, sql=sql, params=params) except Exception as e: log.error('S3 stage \'{}\' creation 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 create_file_format(session, format_type, db, schema): """Create Snowflake file format.""" if format_type == 'ingest': sql = CREATE_S3_CSV_INPUT_FILE_FORMAT.format( snowflake_database=db, snowflake_schema=schema) elif format_type == 'output': sql = CREATE_S3_CSV_OUTPUT_FILE_FORMAT.format( snowflake_database=db, snowflake_schema=schema) else: raise ValueError('\'Type\' must be either \'input\' or \'output\'') try: orch_db.run_query(session, sql=sql) except Exception as e: log.error('Create File Format failed.') raise e def get_logged_release_ids(session_id): """Get a list release_ids which were returned by VAPI. Args: session_id: The session_id of the logs to parse Returns: (list) a list release_ids - [{'release_id': 123},] """ params = { 'source_table': db_const.RELEASE_LOG_TABLE.upper(), 'session_id': session_id } results = orch_db.get_rds_mysql_results( sql=SELECT_RELEASE_ID_FROM_LOG, params=params ) res = [x.release_id for x in results] return res def get_logged_product_and_project(session_id): """Get a list of product_codes and project_ids's from the ingestion log. Args: session_id: (str) The session_id of the logs to parse Returns: (list) a list of dicts filled with 'project_id', 'product_codes' """ params = { 'ingestion_log_table': db_const.RELEASE_LOG_TABLE.upper(), 'session_id': session_id } results = orch_db.get_rds_mysql_results( sql=SELECT_PRODUCT_CODE_FROM_SESSION_ID, params=params ) return results @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_sme_project_names_from_product_codes(session, product_code_list): """Get a list of project names from a list of upcs.""" params = { 'product_code_list': '\'{}\''.format('\', \''.join( map(str, product_code_list))), 'snowflake_db': db_const.SNOWFLAKE_FACTS_DB, 'snowflake_schema': config.ENVIRONMENT, 'snowflake_table': db_const.SNOWFLAKE_LABELCOPY_STAGING, 'snowflake_sony_db': db_const.SNOWFLAKE_SONY_DB, 'snowflake_sony_schema': db_const.SNOWFLAKE_SONY_SCHEMA, 'snowflake_sony_table': db_const.SNOWFLAKE_SONY_PRODUCT_MASTER_TABLE } results = orch_db.run_query( session, sql=SELECT_PROJECT_NAME_FROM_SME_LABELCOPY_BY_UPC, params=params) return results