import time from datetime import datetime, timedelta from typing import Any, Generator from apollo_main_db.apple.models import AppleMusicContainerStreams, AppleMusicContainerStreamSummary, ApplePlaylist, \ AppleWeeklyTopPlaylist from smelog.factory import BoundLogger from sqlalchemy import and_, func, or_ from apple_top_playlists_aggregation.const import DISTRIBUTOR_IDS, TOP_PLAY_LISTS_NUMBER from apple_top_playlists_aggregation.main_db import session_scope def date_range(start_date: datetime.date, end_date: datetime.date) -> Generator[datetime.date, None, None]: """Generator for dates between two dates Args: start_date (datetime.date): Start date. end_date (datetime.date): End date. Returns: Generator[datetime.date]: Dates. """ for n in range(int((start_date - end_date).days) + 1): yield start_date - timedelta(n) def handler(logger: BoundLogger) -> Any: """Aggregation job for apple top 250 playlist Job fetches data from AppleMusicContainerStreamSummary, sorts it and write to AppleWeeklyTopPlaylist Job runs every hour. Args: logger: Logger instance Returns: JSON serializable response """ start_time = time.time() with session_scope() as session: country_codes = session.query( AppleMusicContainerStreamSummary.country_code).distinct( AppleMusicContainerStreamSummary.country_code).all() # AppleMusicContainerStreams contains data for yesterday last_date = datetime.today().date() - timedelta(days=1) dates = [ and_(AppleMusicContainerStreams.date == date) for date in date_range(last_date, last_date - timedelta(days=6))] with session_scope() as session: for country_code, in country_codes: existing_top_playlists_query = session.query( AppleWeeklyTopPlaylist.playlist_id ).filter( AppleWeeklyTopPlaylist.country_code == country_code, AppleWeeklyTopPlaylist.date == last_date.strftime("%Y-%m-%d"), ) new_top_playlists_query = session.query( AppleMusicContainerStreams.container_id ).join( ApplePlaylist, ApplePlaylist.id == AppleMusicContainerStreams.container_id ).filter( AppleMusicContainerStreams.distributor_id.in_(DISTRIBUTOR_IDS), or_(*dates) ).group_by( AppleMusicContainerStreams.container_id ).order_by( func.sum(AppleMusicContainerStreams.streams).desc(), AppleMusicContainerStreams.container_id ) if country_code != 'global': new_top_playlists_query = new_top_playlists_query.filter( AppleMusicContainerStreams.country_code == country_code) existing_top_playlists = existing_top_playlists_query.order_by( AppleWeeklyTopPlaylist.rank.asc()).all() new_top_playlists = new_top_playlists_query.all()[:TOP_PLAY_LISTS_NUMBER] if not new_top_playlists or existing_top_playlists == new_top_playlists: continue existing_top_playlists_query.delete() session.add_all([ AppleWeeklyTopPlaylist( playlist_id=playlist_id, country_code=country_code, date=last_date.strftime("%Y-%m-%d"), rank=idx ) for idx, (playlist_id,), in enumerate(new_top_playlists, start=1) ]) session.commit() logger.info(f'Added rows for {country_code} {time.time() - start_time} sec') logger.info(f'Done {time.time() - start_time} sec',) return {}