"""Snowflake connector class for Shazam flow.""" import datetime from snowflake_connector.etl_connector import SQLLoader from snowflake_connector.etl_connector import SnowflakeSQLExecutor import config # Load SQL templates sql_loader = SQLLoader(__file__) tracks_sql_loader = SQLLoader(__file__, folder='/tracks_queries') class ShazamStatsSFExecutor(SnowflakeSQLExecutor): """Helper class to abstract Snowflake operations.""" @property def stats_name(self): """Name of the stats. Should match dir name of this stats, stats_name in config.py of a statistics. Returns: str: Stats name, e.g. 'shazam_statistics'. """ return config.STATS_NAME @property def artist_to_track_table(self): """Name of the artist_to_track table for the statistics. Returns: str: {stats}_artist_to_track table name. """ return config.snowflake_table_names['shazam_artist_to_track'] @property def shazam_statistics_table(self): """Name of the artist_to_track table for the statistics. Returns: str: {stats}_artist_to_track table name. """ return config.snowflake_table_names['shazam_statistics'] @property def tracks_statistics_table(self): """Name of the artist_to_track table for the statistics. Returns: str: {stats}_artist_to_track table name. """ return config.snowflake_table_names['tracks_statistics'] @property def staging_raw_table(self): """Name of the staging raw table for the statistics. Returns: str: staging raw table name. """ return config.snowflake_table_names['staging_raw'] @property def tracks_staging_raw_table(self): """Name of the staging raw table for the statistics. Returns: str: staging raw table name. """ return config.snowflake_table_names['tracks_staging_raw'] def create_shazam_statistics(self): """Create shazam statistics historical table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.shazam_statistics_table) return self.fetchone_query(sql_loader, 'create_shazam_statistics', params) def create_shazam_tracks_statistics(self): """Create shazam statistics historical table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.tracks_statistics_table) return self.fetchone_query(tracks_sql_loader, 'create_shazam_tracks_statistics', params) def create_staging_raw_shazam_statistics(self): """Create staging raw shazam statistics historical table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.staging_raw_table) return self.fetchone_query(sql_loader, 'create_shazam_statistics', params) def create_staging_raw_shazam_tracks_statistics(self): """Create staging raw shazam statistics historical table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.tracks_staging_raw_table) return self.fetchone_query(tracks_sql_loader, 'create_shazam_tracks_statistics', params) def select_artists_ids(self): """Selecting artists Shazam IDs from shazam_artist_to_track.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.artist_to_track_table) result = self.fetchall_query(sql_loader, 'select_artists', params) return result def select_tracks_ids(self): """Selecting tracks Shazam IDs.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.tracks_statistics_table) result = self.fetchall_query(tracks_sql_loader, 'select_tracks', params) return result def select_artists_stats(self): """Load temp staging raw table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.shazam_statistics_table) result = self.fetchall_query(sql_loader, 'select_artists_stats', params) return result def select_tracks_stats(self): """Load temp staging raw table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.tracks_statistics_table) result = self.fetchall_query(tracks_sql_loader, 'select_tracks_stats', params) return result def delete_from_staging_raw(self): """Selecting artists links from instagram_artist_to_track.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], staging_raw_table=self.staging_raw_table, date=datetime.datetime.today().strftime("%Y-%m-%d")) result = self.fetchall_query(sql_loader, 'delete_from_staging_raw', params) return result def update_artist_stats(self, stats: dict): """Load main statistics table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.shazam_statistics_table, shazam_id=stats.get('id'), current_shazams=stats.get('current_shazams'), last_week_shazams=stats.get('last_week_shazams'), current_weekly_change_in_shazams=stats.get('current_weekly_change_in_shazams'), weekly_change_in_shazams=stats.get('weekly_change_in_shazams'), daily_change_in_shazams=stats.get('daily_change_in_shazams'), percentage_change_in_shazams=stats.get('percentage_change_in_shazams'), acceleration_shazams=stats.get('acceleration_shazams'), average_shazams=stats.get('average_shazams'), current_tracks_count=stats.get('current_tracks_count'), last_week_tracks_count=stats.get('last_week_tracks_count'), weekly_change_in_tracks_count=stats.get('weekly_change_in_tracks_count'), percentage_change_in_tracks_count=stats.get('percentage_change_in_tracks_count'), shazams_per_track=stats.get('shazams_per_track'), last_processing_date=stats.get('last_processing_date').strftime("%Y-%m-%d"), last_week_processing_date=stats.get('last_week_processing_date').strftime("%Y-%m-%d") ) result = self.fetchall_query(sql_loader, 'update_artists_stats', params) return result def update_artist_stats_last_processing_date(self, artist: str, last_processing_date: datetime): """Update artists statistics table on failure.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.shazam_statistics_table, last_processing_date=last_processing_date, artist=artist ) result = self.fetchall_query(sql_loader, 'update_artists_stats_last_processing_date', params) return result def select_weekly_change(self, artist, first_day_of_period=0, last_day_of_period=-7): """Select artist weekly change in metrics from staging raw table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.staging_raw_table, artist=artist, first_day_of_period=first_day_of_period, last_day_of_period=last_day_of_period ) result = self.fetchall_query(sql_loader, 'select_weekly_change', params) try: weeklies = result[0] if len(result[0]) == 2 else [0] * 2 except IndexError: weeklies = [0] * 2 return dict(zip(['shazams', 'tracks_count'], list(weeklies))) def select_track_weekly_change(self, track_id, first_day_of_period=0, last_day_of_period=-7): """Select track weekly change in metrics from staging raw table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.tracks_staging_raw_table, track_id=track_id, first_day_of_period=first_day_of_period, last_day_of_period=last_day_of_period ) result = self.fetchall_query(tracks_sql_loader, 'select_weekly_change', params) try: weekly = result[0][0] if len(result[0]) == 1 else 0 except IndexError: weekly = 0 return dict(shazams=weekly) def update_stating_raw_shazam_stats(self): """Load staging raw table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.shazam_statistics_table, staging_raw_table_name=self.staging_raw_table ) result = self.fetchall_query(sql_loader, 'update_staging_raw_shazam_stats', params) return result def insert_artist_stats(self, artist, stats, table_name=config.snowflake_table_names['shazam_statistics']): """First Shazam statistics insertion.""" stats['the_latest_release'] = '2000-01-01' if stats.get('the_latest_release') \ is None else stats.get('the_latest_release') params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=table_name, artist=artist, shazam_id=stats.get('id'), current_shazams=stats.get('current_shazams'), last_week_shazams=stats.get('last_week_shazams'), current_weekly_change_in_shazams=stats.get('current_weekly_change_in_shazams'), weekly_change_in_shazams=stats.get('weekly_change_in_shazams'), daily_change_in_shazams=stats.get('daily_change_in_shazams'), percentage_change_in_shazams=stats.get('percentage_change_in_shazams'), acceleration_shazams=stats.get('acceleration_shazams'), average_shazams=stats.get('average_shazams'), current_tracks_count=stats.get('current_tracks_count'), last_week_tracks_count=stats.get('last_week_tracks_count'), weekly_change_in_tracks_count=stats.get('weekly_change_in_tracks_count'), percentage_change_in_tracks_count=stats.get('percentage_change_in_tracks_count'), shazams_per_track=stats.get('shazams_per_track'), last_processing_date=stats.get('last_processing_date').strftime("%Y-%m-%d"), last_week_processing_date=stats.get('last_week_processing_date').strftime("%Y-%m-%d") ) result = self.fetchall_query(sql_loader, 'insert_artists_stats', params) return result def update_tracks_stats(self, track_id, stats): """Load main statistics table.""" stats = stats[track_id] params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.tracks_statistics_table, track_id=track_id, track_title=stats.get('title'), artist=stats.get('artist'), shazam_id=stats.get('shazam_id'), track_shazams=stats.get('current_shazams'), last_week_shazams=stats.get('last_week_shazams'), weekly_change_in_shazams=stats.get('weekly_change_in_shazams'), percentage_change_in_shazams=stats.get('percentage_change_in_shazams'), last_processing_date=stats.get('last_processing_date').strftime("%Y-%m-%d"), last_week_processing_date=stats.get('last_week_processing_date').strftime("%Y-%m-%d") ) result = self.fetchall_query(tracks_sql_loader, 'update_shazam_tracks_stats', params) return result def update_staging_raw_shazam_tracks_statistics(self): """Load staging raw table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.tracks_statistics_table, staging_raw_table_name=self.tracks_staging_raw_table ) result = self.fetchall_query(tracks_sql_loader, 'update_staging_raw_shazam_tracks_stats', params) return result def insert_tracks_stats(self, track_id, stats): """First Shazam statistics insertion.""" stats = stats[track_id] params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.tracks_statistics_table, track_id=track_id, track_title=stats.get('title'), artist=stats.get('artist'), shazam_id=stats.get('shazam_id'), track_shazams=stats.get('current_shazams'), last_week_shazams=stats.get('last_week_shazams'), weekly_change_in_shazams=stats.get('weekly_change_in_shazams'), percentage_change_in_shazams=stats.get('percentage_change_in_shazams'), last_processing_date=stats.get('last_processing_date').strftime("%Y-%m-%d"), last_week_processing_date=stats.get('last_week_processing_date').strftime("%Y-%m-%d") ) result = self.fetchall_query(tracks_sql_loader, 'insert_shazam_tracks_stats', params) return result