"""Snowflake connector class for MRC tasks.""" import csv import datetime import config from snowflake_connector.etl_connector import SQLLoader from snowflake_connector.etl_connector import SnowflakeSQLExecutor # Load SQL templates sql_loader = SQLLoader(__file__) looker_sql_loader = SQLLoader(__file__, folder='/looker_queries') class InstagramStatsSFExecutor(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. 'instagram_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['artist_to_track'] @property def instagram_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['instagram_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'] def create_artist_link_table(self): """Create artists links table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.artist_to_track_table) return self.fetchone_query( sql_loader, 'create_artist_to_track', params) def create_instagram_statistics(self): """Create instagram statistics historical table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.instagram_statistics_table) return self.fetchone_query( sql_loader, 'create_instagram_statistics', params) def create_staging_raw_instagram_statistics(self): """Create staging raw instagram 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_instagram_statistics', params) def update_artist_link_table( self, google_sheet_data: list): """Update artists links table. Args: google_sheet_data (list): artists links from Google Sheet. """ for artist in google_sheet_data: params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.artist_to_track_table, artist=artist.get('artist'), instagram_link=artist.get('instagram'), instagram_id=artist.get('id'), spotify_link=artist.get('spotify'), youtube_link=artist.get('youtube'), soundcloud_link=artist.get('soundcloud'), shazam_link=artist.get('shazam'), tiktok_link=artist.get('tiktok') ) self.execute_query( sql_loader, 'update_artist_to_track_from_gs', params) def insert_artist_link_table( self, google_sheet_data: list): """Insert new artists links table. Args: google_sheet_data (list): artists links from Google Sheet. """ for artist in google_sheet_data: params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.artist_to_track_table, artist=artist.get('artist'), instagram_link=artist.get('instagram'), instagram_id=artist.get('id'), spotify_link=artist.get('spotify'), youtube_link=artist.get('youtube'), soundcloud_link=artist.get('soundcloud'), shazam_link=artist.get('shazam'), tiktok_link=artist.get('tiktok'), date_added=artist.get('date') ) self.execute_query( sql_loader, 'insert_artist_to_track_from_gs', params) def delete_artist_link_table(self, artist): """Delete artists links from table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.artist_to_track_table, artist=artist ) self.execute_query(sql_loader, 'delete_artist_to_track', params) return self def delete_from_staging_raw(self): """Select artists links from 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 select_artists_links(self): """Select artists links from 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_artists_stats(self): """Load temp staging raw table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.instagram_statistics_table) result = self.fetchall_query( sql_loader, 'select_artists_stats', 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]) == 5 else [0] * 5 except IndexError: weeklies = [0] * 5 return dict(zip(['followers', 'likes', 'comments', 'posts', 'engagement'], list(weeklies))) def update_artist_stats(self, artist: str, stats: dict): """Load main statistics table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.instagram_statistics_table, link=stats.get('username'), is_verified=stats.get('verified'), current_followers=stats.get('current_followers'), last_week_followers=stats.get('last_week_followers'), weekly_change_in_followers=stats.get('weekly_change_in_followers'), current_weekly_change_in_followers=stats.get( 'current_weekly_change_in_followers'), daily_change_in_followers=stats.get('daily_change_in_followers'), percentage_change_in_followers=stats.get( 'percentage_change_in_followers'), acceleration_followers=stats.get('acceleration_followers'), current_followees=stats.get('current_followees'), current_posts=stats.get('current_posts'), last_week_posts=stats.get('last_week_posts'), weekly_change_in_posts=stats.get('weekly_change_in_posts'), current_likes=stats.get('current_likes'), last_week_likes=stats.get('last_week_likes'), weekly_change_in_likes=stats.get('weekly_change_in_likes'), current_weekly_change_in_likes=stats.get( 'current_weekly_change_in_likes'), daily_change_in_likes=stats.get('daily_change_in_likes'), percentage_change_in_likes=stats.get('percentage_change_in_likes'), acceleration_likes=stats.get('acceleration_likes'), current_comments=stats.get('current_comments'), weekly_change_in_comments=stats.get('weekly_change_in_comments'), current_weekly_change_in_comments=stats.get( 'current_weekly_change_in_comments'), daily_change_in_comments=stats.get('daily_change_in_comments'), percentage_change_in_comments=stats.get( 'percentage_change_in_comments'), engagement=stats.get('current_engagement'), last_week_engagement=stats.get('last_week_engagement'), percentage_change_in_engagement=stats.get( 'percentage_change_in_engagement'), last_processing_date=stats.get( 'last_processing_date').strftime('%d-%m-%Y'), last_week_processing_date=stats.get( 'last_week_processing_date').strftime('%d-%m-%Y'), artist=artist ) 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.instagram_statistics_table, last_processing_date=last_processing_date.strftime('%d-%m-%Y'), artist=artist ) result = self.fetchall_query( sql_loader, 'update_artists_stats_last_processing_date', params) return result def update_stating_raw_instagram_stats(self): """Load staging raw table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=self.instagram_statistics_table, staging_raw_table_name=self.staging_raw_table ) result = self.fetchall_query( sql_loader, 'update_stating_raw_instagram_stats', params) return result def insert_artist_stats( self, artist, stats, table_name=config.snowflake_table_names['instagram_statistics']): """First Instagram statistics insertion.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=table_name, artist=artist, artist_id=stats.get('username'), is_verified=stats.get('verified'), location='', current_followers=stats.get('current_followers'), last_week_followers=stats.get('last_week_followers'), weekly_change_in_followers=stats.get('weekly_change_in_followers'), current_weekly_change_in_followers=stats.get( 'current_weekly_change_in_followers'), daily_change_in_followers=stats.get('daily_change_in_followers'), percentage_change_in_followers=stats.get( 'percentage_change_in_followers'), acceleration_followers=stats.get('acceleration_followers'), current_followees=stats.get('followings'), current_posts=stats.get('current_posts'), last_week_posts=stats.get('last_week_posts'), weekly_change_in_posts=stats.get('weekly_change_in_posts'), current_likes=stats.get('current_likes'), last_week_likes=stats.get('last_week_likes'), weekly_change_in_likes=stats.get('weekly_change_in_likes'), current_weekly_change_in_likes=stats.get( 'current_weekly_change_in_likes'), daily_change_in_likes=stats.get('daily_change_in_likes'), percentage_change_in_likes=stats.get('percentage_change_in_likes'), acceleration_likes=stats.get('acceleration_likes'), current_comments=stats.get('current_comments'), weekly_change_in_comments=stats.get('weekly_change_in_comments'), current_weekly_change_in_comments=stats.get( 'current_weekly_change_in_comments'), daily_change_in_comments=stats.get('daily_change_in_comments'), percentage_change_in_comments=stats.get( 'percentage_change_in_comments'), engagement=stats.get('current_engagement'), last_week_engagement=stats.get('last_week_engagement'), percentage_change_in_engagement=stats.get( 'percentage_change_in_engagement'), last_processing_date=stats.get( 'last_processing_date').strftime('%d-%m-%Y'), last_week_processing_date=stats.get( 'last_week_processing_date').strftime('%d-%m-%Y'), ) result = self.fetchall_query(sql_loader, 'insert_artists_stats', params) return result def execute_looker_query(self, artist, date=datetime.date.today().strftime('%d-%m-%Y')): """Load staging raw table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], staging_raw_table_name=self.staging_raw_table, artist=artist ) result = self.fetchall_query(looker_sql_loader, 'engagement_dashboard_query', params) filename = 'looker_queries_data/' \ '{artist}_engagement_up_to_{date}.csv'. \ format(date=date, artist=artist.replace(' ', '_')) with open(filename, 'w') as file: writer = csv.writer(file) writer.writerow(['date', 'engagement']) writer.writerows(result) return result, filename