import logging from datetime import datetime from typing import Dict from sqlalchemy.dialects.postgresql import insert from sqlalchemy.engine import create_engine from backfill_lambda.config import Config from backfill_lambda.models import Playlist from backfill_lambda.queries import ( insert_dapd_apple_music_storefront, insert_dapd_spotify_playlists, insert_dapd_spotify_playlists_countries, insert_dapd_spotify_storefront, select_dapd_apple_music_markets, select_dapd_apple_playlists, select_dapd_apple_playlists_2, select_dapd_apple_playlists_count, select_dapd_apple_playlists_count_2, select_dapd_spotify_markets, select_dapd_spotify_playlists, select_dapd_spotify_playlists_count, select_dapd_spotify_playlists_countries, select_dapd_spotify_playlists_countries_count, ) # logging.basicConfig( # level=logging.INFO, # format='%(asctime)s %(message)s') logging.getLogger().setLevel(logging.INFO) logger = logging.getLogger('jobs') def dapd_spotify_playlists_backfill(event, config: Config) -> Dict[str, int]: now = datetime.now().strftime('%Y-%m-%d %H:%M:%S') extension = { 'data_source_id': 2, 'created_at': now, 'label': '', 'expired_at': now, } source_engine = create_engine(config.apollo_uri).execution_options(streaming_results=True) dest_engine = create_engine(config.workflow_uri) offset = event.get('offset', 0) logger.info('Task started') with source_engine.connect() as source_conn: logger.info('Source connection estabilished') count = source_conn.execute(select_dapd_spotify_playlists_count) count = count.fetchone()[0] logger.info('Total playlists to process: %s', count) while offset < count: result = source_conn.execute(select_dapd_spotify_playlists, {'offset': offset}) while True: chunk = result.fetchmany(config.db_read_chunk_size) if not chunk: break with dest_engine.connect() as dest_conn: with dest_conn.begin(): markets = dest_conn.execute(select_dapd_spotify_markets) markets = markets.fetchall() markets = [market[0] for market in markets] for row in chunk: if row['storefront_id'] not in markets: dest_conn.execute( insert_dapd_spotify_storefront, {'market': row['storefront_id']}, ) row = prepare_row(row, extension) dest_conn.execute(insert_dapd_spotify_playlists, row) offset += 50000 logger.info('Processed %s playlists', offset) return {'status_code': 200, 'processed': count} def dapd_spotify_playlists_countries_backfill(event, config: Config) -> Dict[str, int]: extension = {} source_engine = create_engine(config.apollo_uri).execution_options(streaming_results=True) dest_engine = create_engine(config.workflow_uri) offset = event.get('offset', 0) logger.info('Task started') with source_engine.connect() as source_conn: logger.info('Source connection estabilished') count = source_conn.execute(select_dapd_spotify_playlists_countries_count) count = count.fetchone()[0] logger.info('Total playlists to process: %s', count) while offset < count: result = source_conn.execute(select_dapd_spotify_playlists_countries, {'offset': offset}) i = offset while True: chunk = result.fetchmany(config.db_read_chunk_size) if not chunk: break with dest_engine.connect() as dest_conn: for row in range(len(chunk)): chunk[row] = tuple(prepare_row(chunk[row], extension)) query = insert_dapd_spotify_playlists_countries.format(str(chunk)[1:-1]) dest_conn.execute(query) i += config.db_read_chunk_size logger.info('Processed %i playlists', i) offset += 50000 logger.info('Finished') return {'status_code': 200, 'processed': count} def dapd_apple_playlists_backfill(event, config: Config): storefronts = { "top_storefronts": [ "us", "gb", "au", "ca", "de", "fr" ], "major_storefronts": [ "us", "gb", "cn", "au", "ca", "ru", "de", "fr", "mx", "br", "in", "za", "kr", "it", "tw", "tr", "ch", "es", "th", "hk", "ae", "nz", "dk", "id", "il", "nl", "se", "no" ], } now = datetime.now().strftime('%Y-%m-%d %H:%M:%S') extension = { 'data_source_id': 1, 'created_at': now, 'label': '', 'expired_at': now, } source_engine = create_engine(config.apollo_uri).execution_options(streaming_results=True) dest_engine = create_engine(config.workflow_uri) offset = event.get('offset', 0) logger.info('Task started') with source_engine.connect() as source_conn: logger.info('Source connection estabilished') count = source_conn.execute(select_dapd_apple_playlists_count) count = count.fetchone()[0] logger.info('Total playlists to process: %s', count) while offset < count: result = source_conn.execute( select_dapd_apple_playlists, {'offset': offset, **storefronts}, ) c = 0 logger.info('Select query executed') while True: chunk = result.fetchmany(config.db_read_chunk_size) logger.info('Chunk fetched') if not chunk: break c += len(chunk) query = insert(Playlist).values([{**playlist, **extension} for playlist in chunk]) query = query.on_conflict_do_nothing( index_elements=['id', 'data_source_id', 'storefront_id'], ) with dest_engine.connect() as dest_conn: markets = dest_conn.execute(select_dapd_apple_music_markets) markets = markets.fetchall() markets = [market[0] for market in markets] for playlist in chunk: if playlist.storefront_id not in markets: dest_conn.execute( insert_dapd_apple_music_storefront, {'market': playlist.storefront_id},) dest_conn.execute(query) logger.info('Chunk updated') offset += 10000 logger.info('Processed %s playlists', offset) return {'status_code': 200, 'processed': offset} def dapd_apple_playlists_backfill_2(event, config: Config): storefronts = { "top_storefronts": [ "us", "gb", "au", "ca", "de", "fr" ], "major_storefronts": [ "us", "gb", "cn", "au", "ca", "ru", "de", "fr", "mx", "br", "in", "za", "kr", "it", "tw", "tr", "ch", "es", "th", "hk", "ae", "nz", "dk", "id", "il", "nl", "se", "no" ], } now = datetime.now().strftime('%Y-%m-%d %H:%M:%S') extension = { 'data_source_id': 1, 'created_at': now, 'label': '', 'expired_at': now, } source_engine = create_engine(config.apollo_uri).execution_options(streaming_results=True) dest_engine = create_engine(config.workflow_uri) offset = event.get('offset', 0) logger.info('Task started') with source_engine.connect() as source_conn: logger.info('Source connection estabilished') count = source_conn.execute(select_dapd_apple_playlists_count_2) count = count.fetchone()[0] logger.info('Total playlists to process: %s', count) while offset < count: result = source_conn.execute( select_dapd_apple_playlists_2, {'offset': offset, **storefronts}, ) c = 0 logger.info('Select query executed') while True: chunk = result.fetchmany(config.db_read_chunk_size) logger.info('Chunk fetched') if not chunk: break c += len(chunk) query = insert(Playlist).values([{**playlist, **extension} for playlist in chunk]) query = query.on_conflict_do_nothing( index_elements=['id', 'data_source_id', 'storefront_id'], ) with dest_engine.connect() as dest_conn: markets = dest_conn.execute(select_dapd_apple_music_markets) markets = markets.fetchall() markets = [market[0] for market in markets] for playlist in chunk: if playlist.storefront_id not in markets: dest_conn.execute( insert_dapd_apple_music_storefront, {'market': playlist.storefront_id},) dest_conn.execute(query) logger.info('Chunk updated') offset += 10000 logger.info('Processed %s playlists', offset) return {'status_code': 200, 'processed': offset} def spotify_playlists_backfill(event, config: Config): raise Exception('Not implemented') def apple_playlists_backfill(event, config: Config): raise Exception('Not implemented') def prepare_row(row, extension): row = {**row, **extension} if row.get('storefront_id') in (None, 'null', ''): row['storefront_id'] = 'not_defined' return row if row.get('storefront_id') == '_gl': row['storefront_id'] = 'global' return row row['storefront_id'] = row.get('storefront_id').lower() return row