"""Shazam scrapping flow.""" import logging import asyncio import datetime from shazamio.client import HTTPClient from aiohttp import ClientConnectorError from shazamio import Shazam, serialize_track from shazamio.exceptions import FailedDecodeJson from shazamio.models import Request from helpers import catch_all_and_print from snowflake_executor import ShazamStatsSFExecutor import config def today(): """Return today's date.""" return datetime.date.today() def zero_if_none(value): return value if value is not None else 0 def get_id_from_link(text: str) -> dict: # text = config.artists_to_track.get(artist) """Get ID from Shazam profile link.""" path = text.rsplit('/') if path[-1].isnumeric(): return dict(id=int(''.join(c for c in text.rsplit('/')[-1] if c.isnumeric())), type='new') else: logging.info(text + ' seems to be an old-type link.') return dict(id=int(''.join(c for c in text.rsplit('/')[-2] if c.isnumeric())), type='old') def if_first_time(stat, first_time): """Return stat depending onn first_time bool value.""" return stat if not first_time else 0 async def get_track_id(adam_id): track_info = await HTTPClient.request( 'GET', config.TRACK_INFO.format(adam_id), headers=Request.HEADERS) track_id = track_info['data'][0].get('id') return dict(name=track_info['resources']['shazam-songs'] [track_id]['attributes'].get('title'), id=track_id) class ShazamTracker: """Shazam Tracker class for scrapping data.""" def __init__(self): logging.basicConfig(level=logging.INFO) self.__shazam__ = Shazam() self.__executor__ = ShazamStatsSFExecutor(sf_config=config.SF_CONFIG) self.__loop__ = asyncio.get_event_loop() self.__profiles__ = dict() self.__tracks__ = dict() self.__last_week_stats__ = dict() self.__last_week_tracks_stats__ = dict() self.__first_time__ = False def get_shazam(self): """Get Shazam API instance.""" return self.__shazam__ def get_executor(self): """Get SnowFlake executor instance.""" return self.__executor__ def get_loop(self): """Get asyncio event loop.""" return self.__loop__ def get_profiles(self): """Get profiles.""" return self.__profiles__ def get_tracks(self): """Get tracks.""" return self.__tracks__ def get_last_week_stats(self, tracks=False): """Get last week stats.""" if tracks: return self.__last_week_tracks_stats__ return self.__last_week_stats__ def set_last_week_stats(self): self.__last_week_stats__.update( dict((artist_name, {'shazam_id': shazam_id, 'current_shazams': current_shazams, 'last_week_shazams': last_week_shazams, 'weekly_change_in_shazams': weekly_change_in_shazams, 'percentage_change_in_shazams': percentage_change_in_shazams, 'acceleration_shazams': acceleration_shazams, 'average_shazams': average_shazams, 'current_tracks_count': current_tracks_count, 'last_week_tracks_count': last_week_tracks_count, 'weekly_change_in_tracks_count': weekly_change_in_tracks_count, 'percentage_change_in_tracks_count': percentage_change_in_tracks_count, 'shazams_per_track': shazams_per_track, 'last_processing_date': last_processing_date, 'last_week_processing_date': last_week_processing_date}) for artist_name, shazam_id, current_shazams, last_week_shazams, weekly_change_in_shazams, percentage_change_in_shazams, acceleration_shazams, average_shazams, current_tracks_count, last_week_tracks_count, weekly_change_in_tracks_count, percentage_change_in_tracks_count, shazams_per_track, last_processing_date, last_week_processing_date in self.__executor__.select_artists_stats())) def set_last_week_tracks_stats(self): self.__last_week_tracks_stats__.update( dict((track_id, {'track_title': track_title, 'artist': artist, 'shazam_id': shazam_id, 'current_shazams': track_shazams, 'last_week_shazams': last_week_shazams, 'weekly_change_in_shazams': weekly_change_in_shazams, 'percentage_change_in_shazams': percentage_change_in_shazams, 'last_processing_date': last_processing_date, 'last_week_processing_date': last_week_processing_date}) for track_title, track_id, artist, shazam_id, track_shazams, last_week_shazams, weekly_change_in_shazams, percentage_change_in_shazams, last_processing_date, last_week_processing_date in self.__executor__.select_tracks_stats())) def create_tables(self): """Create tables if not exist and read last week stats.""" exists = ' table already exists.' executor = self.get_executor() self.__profiles__.update(dict((artist, dict(id=get_id_from_link(link)['id'], type=get_id_from_link(link)['type'])) for artist, link in executor.select_artists_ids() if link is not None)) query_result = executor.create_shazam_statistics() if 'already exists' in query_result[0]: logging.info('%s%s', config.snowflake_table_names['shazam_statistics'].upper(), exists) self.set_last_week_stats() else: self.__first_time__ = True query_result = executor.create_staging_raw_shazam_statistics() if 'already exists' in query_result[0]: logging.info('%s%s\n', config.snowflake_table_names['staging_raw'].upper(), exists) q1_, q2_ = executor.create_shazam_tracks_statistics(), \ executor.create_staging_raw_shazam_tracks_statistics() for q_ in [q1_, q2_]: if 'already exists' in q_[0]: logging.info(q_[0]) self.__tracks__.update(dict((track_id, dict(title=track_title)) for track_title, track_id in executor.select_tracks_ids())) self.set_last_week_tracks_stats() deleted = executor.delete_from_staging_raw() logging.info('%i rows have been deleted.\n', deleted[0][0]) @catch_all_and_print def compare_current_and_last_week_stats( self, current: dict, stat_name: str, lasts: dict, artist: str, today=today()): """Return change in stats comparing current ones with last week ones.""" last = dict() try: last = {artist: dict(zip( lasts.get(artist).keys(), map(lambda key: 0 if lasts.get(artist).get(key) is None else lasts.get(artist).get(key), lasts.get(artist).keys())))} except AttributeError: last[artist] = { 'shazam_id': None, 'current_shazams': 0, 'last_week_shazams': 0, 'weekly_change_in_shazams': 0, 'percentage_change_in_shazams': 0, 'acceleration_shazams': 0, 'average_shazams': 0, 'current_tracks_count': 0, 'last_week_tracks_count': 0, 'weekly_change_in_tracks_count': 0, 'percentage_change_in_tracks_count': 0, 'shazams_per_track': 0, 'last_processing_date': today, 'last_week_processing_date': today } if not artist.isnumeric() else { 'track_title': None, 'artist': None, 'shazam_id': None, 'current_shazams': 0, 'last_week_shazams': 0, 'weekly_change_in_shazams': 0, 'percentage_change_in_shazams': 0, 'last_processing_date': today, 'last_week_processing_date': today} last[artist]['current_' + stat_name] = current[stat_name] last[artist]['last_week_' + stat_name] = current[stat_name] self.__first_time__ = True previous_stat = float(last[artist]. \ get('current_' + stat_name)) last_week_stat = float(last[artist]. \ get('last_week_' + stat_name)) last_week_processing_date = last[artist]. \ get('last_week_processing_date') current = previous_stat if current[stat_name] is None else current[stat_name] daily_change = current - previous_stat logging.info('Current count of %s is %i', stat_name.replace('_', ' '), current) logging.info('Last week the count of %s was %i', stat_name, last_week_stat) if artist.isnumeric(): current_weekly_change = self.get_executor().select_track_weekly_change(int(artist))[ stat_name] + daily_change last_weekly_change = self.get_executor().select_track_weekly_change( int(artist), first_day_of_period=-6, last_day_of_period=-14)[ stat_name] else: current_weekly_change = self.get_executor().select_weekly_change(artist)[ stat_name] + daily_change last_weekly_change = self.get_executor().select_weekly_change( artist, first_day_of_period=-6, last_day_of_period=-14)[ stat_name] # last_weekly_change also calculating, but is not pushed to DB for now percentage_change = (current_weekly_change - last_weekly_change) / last_weekly_change \ if (last_weekly_change) != 0 else 0.0 acceleration = (current_weekly_change - last_weekly_change) \ / last_weekly_change if last_weekly_change != 0 else 0 logging.info('Change in %s is %i which is %s%s', stat_name, current_weekly_change, str(percentage_change), '%') logging.info(stat_name[0].upper() + stat_name[1:].replace('_', ' ') + ' acceleration is ' + str(acceleration)) return {'current_' + stat_name: current, 'last_week_' + stat_name: current - current_weekly_change, 'daily_change_in_' + stat_name: daily_change, 'current_weekly_change_in_' + stat_name: current_weekly_change, # current_weekly_change as default 'weekly_change_in_' + stat_name: last_weekly_change, # last_weekly_change as default 'percentage_change_in_' + stat_name: percentage_change, 'acceleration_' + stat_name: acceleration, 'last_processing_date': today, 'last_week_processing_date': today if today >= last_week_processing_date + \ datetime.timedelta(days=7) else last_week_processing_date } @catch_all_and_print async def get_artist_top_tracks(self, artist_id: int, shazam_version: str) -> list: """Return a list of artist's all tracks. def artist_top_tracks() has default songs limit as 200. Currently, it's enough for tracking all tracks stats.""" if shazam_version.startswith('new'): top_artist_tracks = await HTTPClient.request( 'GET', config.ARTIST_TOP_TRACKS_NEW.format(artist_id), headers=Request.HEADERS) return [dict(name=t['attributes']['name'], id=t['id']) for t in top_artist_tracks['data']] else: top_artist_tracks = await self.get_shazam().artist_top_tracks( artist_id=artist_id) return [dict(name=serialize_track(data=t).title, id=serialize_track(data=t).key) for t in top_artist_tracks['tracks']] @catch_all_and_print async def get_listenings_count(self, track_id: int) -> int: """Return Shazams by Shazam track ID.""" count = await self.get_shazam().listening_counter(track_id=track_id) return count.get('total') def get_tracks_(self, artist_name, log=False): """Loading artist's Shazam tracks IDs. Args: artist_name (str): artist's name. log (bool): log loading info or not. """ artist_id, shazam_version = self.get_profiles()[artist_name].get('id'), \ self.get_profiles()[artist_name].get('type') tracks = self.get_loop().run_until_complete( self.get_artist_top_tracks(artist_id, shazam_version)) if log: logging.info(artist_name + ', Shazam ID: ' + str(artist_id)) return tracks, artist_id def get_common_stats(self, artist): """Return count of followers and popularity. Args: artist (str): artist's name.""" profiles = self.get_profiles() stats = self.get_last_week_stats() tracks_to_table = self.get_tracks() tracks_stats = self.get_last_week_stats(tracks=True) tracks = self.get_tracks_(artist, log=True)[0] # for track in tracks: # track.update(shazams=self.get_loop().run_until_complete( # self.get_listenings_count(track_id=track.get('id')))) for track in tracks: if self.get_profiles()[artist].get('type') is 'new': track_id = self.get_loop().run_until_complete( get_track_id(track.get('id'))) track.update(track_id=track_id.get('id')) else: track_id = track track.update(shazams=self.get_loop().run_until_complete( self.get_listenings_count(track_id=track_id.get('id')))) for track in tracks: id_ = str(track.get('id')) try: tracks_to_table[id_].update(title=track.get('name')) except KeyError: tracks_to_table.update({id_: {'title': track.get('name')}}) tracks_to_table[id_].update(artist=artist) tracks_to_table[id_].update(shazam_id=profiles[artist].get('id')) tracks_to_table[id_].update(self.compare_current_and_last_week_stats( dict(shazams=track.get('shazams')), 'shazams', tracks_stats, id_)) if id_ not in tracks_stats.keys(): self.get_executor().insert_tracks_stats(id_, tracks_to_table) else: self.get_executor().update_tracks_stats(id_, tracks_to_table) shazams = sum(track['shazams'] for track in tracks) shazams_per_track = ',\n'.join([t.get('name') + ' - ' + str(t.get('shazams')) for t in tracks]) profiles[artist].update(self.compare_current_and_last_week_stats( dict(shazams=shazams), 'shazams', stats, artist)) profiles[artist].update(self.compare_current_and_last_week_stats( dict(tracks_count=len(tracks)), 'tracks_count', stats, artist)) profiles[artist].update(dict( average_shazams=shazams / len(tracks), shazams_per_track=shazams_per_track)) def update_new_artists_stats(self, profile): """Update all stats for artists at SnowFlake.""" if self.__first_time__: self.get_executor().insert_artist_stats(profile, self.get_profiles()[profile]) logging.info("%s's stats has been inserted.\n", profile) self.get_executor().insert_artist_stats(profile, self.get_profiles()[profile], table_name=self.get_executor().staging_raw_table) def update_missing_days_stats(self, artist): if not artist.isnumeric(): missing_days_stats, profiles, stats = \ {}, self.get_profiles(), self.get_last_week_stats() missing_days = (datetime.date.today() - stats[artist]. get('last_processing_date')).days days = [datetime.date.today() - datetime.timedelta(days=i) for i in range(missing_days)] missing_days_stats.update(dict( artist=artist, id=profiles[artist].get('id'), shazams_per_track=profiles[artist].get('shazams_per_track') )), days.sort() if missing_days > 0: missing_daily_shazams = profiles[artist].get( 'daily_change_in_shazams') / missing_days missing_daily_tracks_count = zero_if_none(profiles[artist].get( 'daily_change_in_tracks_count')) / missing_days for day in days: today_, stats = day, \ self.get_last_week_stats() # calculating actual metrics for shazams, tracks count missing_days_stats.update(self.compare_current_and_last_week_stats( dict(shazams=float(stats[artist].get('current_shazams')) + missing_daily_shazams), 'shazams', stats, artist, today=today_)) missing_days_stats.update(self.compare_current_and_last_week_stats( dict(tracks_count=stats[artist].get('current_tracks_count') + missing_daily_tracks_count), 'tracks_count', stats, artist, today=today_)) missing_days_stats.update(dict( average_shazams=(float(stats[artist].get('current_shazams')) + missing_daily_shazams) / (stats[artist].get('current_tracks_count') + missing_daily_tracks_count))) # updating stats table with new calculations to keep the growth gradual self.get_executor().update_artist_stats(missing_days_stats) # updating staging raw table with new calculations to get correct weekly change self.get_executor().insert_artist_stats( artist, missing_days_stats, table_name=config.snowflake_table_names['staging_raw']) # setting last week stats with the newest data for the missing day self.set_last_week_stats() else: self.get_executor().update_artist_stats(self.get_profiles()[artist]) def get_all_stats_and_update(self): """Update stats for every single artist.""" tracking_start = datetime.datetime.now() for artist in self.get_profiles(): self.__first_time__ = False try: start = datetime.datetime.now() self.get_common_stats(artist) finish = datetime.datetime.now() logging.info('For %s it took %i seconds.', artist, (finish - start).total_seconds()) # updates staging raw with missing days stats if artist in self.get_last_week_stats(): self.update_missing_days_stats(artist) self.update_new_artists_stats(artist) except ConnectionError: logging.error('An error occurred on requests side.') except FailedDecodeJson as error: if 'URL is invalid' in error.args[0]: logging.error('Provided Shazam link for %s is invalid.\n', artist) except ClientConnectorError: logging.error("Cannot connect to Shazam for getting %s's stats.", artist) tracking_finish = datetime.datetime.now() logging.info('For all the artists it took %i seconds.\n', (tracking_finish - tracking_start).total_seconds()) def update_staging_raw_table(self): """Update appropriate staging raw table.""" logging.info( 'All the stats has been updated at %s.\n', config.snowflake_table_names['shazam_statistics']) self.get_executor().update_stating_raw_shazam_stats() self.get_executor().update_staging_raw_shazam_tracks_statistics() @catch_all_and_print def perform_tracking(self): """Perform tracking flow.""" self.create_tables() self.get_all_stats_and_update() # self.update_staging_raw_table() self.get_executor().update_staging_raw_shazam_tracks_statistics() if __name__ == '__main__': tracker = ShazamTracker() tracker.perform_tracking()