import pandas as pd import numpy as np import streamlit as st import pytz import datetime as dt from dateutil.relativedelta import relativedelta, FR import snowflake.connector from snowflake.connector.pandas_tools import write_pandas import redshift_connector import os import boto3 from cryptography.hazmat.backends import default_backend from cryptography.hazmat.primitives import serialization try: sf_user = os.environ['SNOWFLAKE_USER'] sf_account = os.environ['SNOWFLAKE_ACCOUNT'] sf_warehouse = os.environ['SNOWFLAKE_WAREHOUSE'] sf_password = os.environ['PEM_KEY_PASSWORD'] sf_pem_key = os.environ['PEM_KEY'] db_user = os.environ['DB_USER'] db_password = os.environ['DB_PASSWORD'] except: pass try: sf_user = st.secrets['SNOWFLAKE_USER'] sf_account = st.secrets['SNOWFLAKE_ACCOUNT'] sf_warehouse = st.secrets['SNOWFLAKE_WAREHOUSE'] sf_password = st.secrets['PEM_KEY_PASSWORD'] sf_pem_key = st.secrets['PEM_KEY'] db_user = st.secrets['DB_USER'] db_password = st.secrets['DB_PASSWORD'] except: pass try: session = boto3.session.Session() client = session.client( service_name='secretsmanager', region_name='us-east-1' ) sf_user = client.get_secret_value(SecretId="dev/awal-ar/SNOWFLAKE_USER")["SecretString"] sf_account = client.get_secret_value(SecretId="dev/awal-ar/SNOWFLAKE_ACCOUNT")["SecretString"] sf_warehouse = client.get_secret_value(SecretId="dev/awal-ar/SNOWFLAKE_WAREHOUSE")["SecretString"] sf_password = client.get_secret_value(SecretId="dev/awal-ar/PEM_KEY_PASSWORD")["SecretString"] sf_pem_key = client.get_secret_value(SecretId="dev/awal-ar/PEM_KEY")["SecretString"] db_user = client.get_secret_value(SecretId="dev/awal-discovery-main/DB_USER")["SecretString"] db_password = client.get_secret_value(SecretId="dev/awal-discovery-main/DB_PASSWORD")["SecretString"] except: pass p_key = serialization.load_pem_private_key( sf_pem_key.encode('utf-8').decode('unicode_escape').encode("utf-8"), password=sf_password.encode('utf-8'), backend=default_backend() ) pkb = p_key.private_bytes( encoding=serialization.Encoding.DER, format=serialization.PrivateFormat.PKCS8, encryption_algorithm=serialization.NoEncryption()) def snowflake_query(sql_query): ctx = snowflake.connector.connect( user=sf_user, private_key=pkb, account=sf_account, warehouse=sf_warehouse ) cs = ctx.cursor(snowflake.connector.DictCursor) try: cs.execute(sql_query) result = cs.fetchall() finally: cs.close() ctx.close() return pd.DataFrame(result) def run_query(query_text): cursor.execute(f"{query_text}") # convert to pd dataframe df_data = cursor.fetchall() col_names = cursor.description col_names = [col[0] for col in col_names] df = pd.DataFrame(df_data,columns=col_names) return df def organize_wtd(df,which_week): df[f'WTD_{which_week}'] = df.groupby(['UNIFIED_SONG_ID','REGION'])['DAILY_STREAMS'].transform('sum') df.drop(['REPORT_DATE','DAILY_STREAMS','WHICH_DAY'],axis=1,inplace=True) df.drop_duplicates(inplace=True,ignore_index=True) df_gl = df.loc[df['REGION']=='Global'].copy() df_gl.rename(columns={f'WTD_{which_week}':f'GL_WTD_{which_week}'},inplace=True) df_gl.drop('REGION',axis=1,inplace=True) df_us = df.loc[df['REGION']=='US'].copy() df_us.rename(columns={f'WTD_{which_week}':f'US_WTD_{which_week}'},inplace=True) df_us.drop('REGION',axis=1,inplace=True) df = df_gl.merge(df_us,how='outer',on='UNIFIED_SONG_ID') return df # organize/pivot; fill NAs def organize_daily_streams(tmp_pivot,which_region): pivot_region = tmp_pivot[['UNIFIED_SONG_ID','WHICH_DAY',which_region]] pivot_region = pd.pivot_table(pivot_region, values=which_region, index=['UNIFIED_SONG_ID'], columns=['WHICH_DAY']).reset_index() pivot_region.fillna(0,inplace=True) # rename columns with prefixes if which_region=='Global': col_prefix = 'GL' elif which_region=='US': col_prefix = 'US' for i in np.arange(7): pivot_region.rename(columns={f'DAY_{i+1}':f'{col_prefix}_DAY_{i+1}'},inplace=True) pivot_region[f'{col_prefix}_DAY_{i+1}'] = pivot_region[f'{col_prefix}_DAY_{i+1}'].astype('int64') return pivot_region # organize/pivot; fill NAs def organize_weekly_streams(tmp_pivot,which_region): pivot_region = tmp_pivot[['UNIFIED_SONG_ID','WHICH_WEEK',which_region]] pivot_region = pd.pivot_table(pivot_region, values=which_region, index=['UNIFIED_SONG_ID'], columns=['WHICH_WEEK']).reset_index() pivot_region.fillna(0,inplace=True) # rename columns with prefixes if which_region=='Global': col_prefix = 'GL' elif which_region=='US': col_prefix = 'US' for i in np.arange(4): pivot_region.rename(columns={f'LW{i+1}':f'{col_prefix}_LW{i+1}'},inplace=True) pivot_region[f'{col_prefix}_LW{i+1}'] = pivot_region[f'{col_prefix}_LW{i+1}'].astype('int64') return pivot_region def pull_metadata_from_spotify(df): # Logic = get spotify track then spotify album to get release date # use left outer join for artist join, because songs with multiple artists will not join in that query # note that we can only pass in x number of ID's per sql call so we might need to split it up isrcs = [] n_calls = np.ceil(len(df)/15000) start_idx = 0 for i in np.arange(n_calls): print(f"Loop {int(i+1)} of {int(n_calls)}") isrc_df = snowflake_query(f""" select distinct t.isrc, t.SPOTIFY_TRACK_ID, t.PREVIEW_URL, a.release_date, artist.genres, artist.followers_latest, a.LABEL from chartmetric.raw_data.spotify t join chartmetric.raw_data.spotify_album a on a.SPOTIFY_ALBUM_ID=t.spotify_album_id left outer join chartmetric.raw_data.spotify_artist artist on artist.SPOTIFY_ARTIST_ID=t.SPOTIFY_ARTIST_ID where 1=1 and t.isrc in {tuple(df['ISRC'].iloc[start_idx:start_idx+15000])} """) if len(isrc_df)>0: isrc_df.drop_duplicates('ISRC',inplace=True) isrcs.append(isrc_df) del isrc_df start_idx = start_idx + 15000 if len(isrcs)>0: isrcs = pd.concat(isrcs,ignore_index=True) else: isrcs = pd.DataFrame() return isrcs def pull_metadata_from_itunes(df): isrc_df = snowflake_query(f""" select t.isrc, t.itunes_track_id, t.PREVIEW_URL, a.release_date, a.genres, concat(a.copyright,', ',a.LABEL) label from chartmetric.raw_data.itunes t join chartmetric.raw_data.itunes_album a on a.itunes_album_id=t.itunes_album_id where 1=1 and t.isrc in {tuple(df['ISRC'])} """) if len(isrc_df)>0: isrc_df.drop_duplicates('ISRC',inplace=True) return isrc_df def pull_metadata_from_amazon(df): isrc_df = snowflake_query(f""" select t.isrc, t.id amazon_track_id, a.release_date, a.genre genres, concat(a.copyright,', ',a.LABEL) label from chartmetric.raw_data.amazon t join chartmetric.raw_data.amazon_album a on a.id=t.AMAZON_ALBUM where 1=1 and t.isrc in {tuple(df['ISRC'])} """) if len(isrc_df)>0: isrc_df.drop_duplicates('ISRC',inplace=True) return isrc_df def pull_metadata_from_deezer(df): isrc_df = snowflake_query(f""" select t.isrc, t.id deezer_track_id, t.PREVIEW_URL, a.release_date, a.LABEL from chartmetric.raw_data.deezer t join chartmetric.raw_data.deezer_album a on a.id=t.DEEZER_ALBUM where 1=1 and t.isrc in {tuple(df['ISRC'])} """) if len(isrc_df)>0: isrc_df.drop_duplicates('ISRC',inplace=True) return isrc_df def upload_df_to_snowflake_table(df,table_name): ctx = snowflake.connector.connect( user=sf_user, private_key=pkb, account=sf_account, warehouse=sf_warehouse, database='awal', schema='awal_ar' ) try: write_pandas(ctx, df, table_name, auto_create_table=True) finally: ctx.close() #################################### #################################### #################################### # connect to DB print('Connecting to DB...') conn = redshift_connector.connect( host='reporting-db.smeanalyticsapps.com', database='smaredshiftdb', port=5439, user=db_user, #insert your Reporting DB username password=db_password #insert your Reporting DB password ) cursor = conn.cursor() # set start time just so we can see how long this takes to run start_time = dt.datetime.now() #################################### print('Setting up dates...') # setting up dates are a little different here bc the table doesn't update regularly. We base it off latest available date. my_query = f""" select max(report_date) from luminate.vw_daily_top_songs_global """ latest_day = run_query(my_query) # today = dt.datetime.now(pytz.timezone('US/Eastern')).date() today = latest_day['max'].iloc[0] + relativedelta(days=2)# - relativedelta(days=1) # print(today) #### set start/end date ranges (friday to the latest building day) ### # current building week if today.weekday() == 5: bw_start_date_tp = today + relativedelta(days=-1, weekday=FR(-2)) else: bw_start_date_tp = today + relativedelta(days=-1, weekday=FR(-1)) bw_end_date_tp = today - relativedelta(days=2) # luminate latest day typically 2 days behind n_building_days = (bw_end_date_tp - bw_start_date_tp).days + 1 # LP building week (last week) bw_start_date_lp_bw = bw_start_date_tp + relativedelta(weeks=-1) bw_end_date_lp_bw = bw_start_date_lp_bw + relativedelta(days=n_building_days-1) # LP (last week) bw_start_date_lp_1 = bw_start_date_tp + relativedelta(weeks=-1) bw_end_date_lp_1 = bw_start_date_lp_1 + relativedelta(days=6) # LP (2 weeks ago) bw_start_date_lp_2 = bw_start_date_lp_1 + relativedelta(weeks=-1) bw_end_date_lp_2 = bw_start_date_lp_2 + relativedelta(days=6) # LP (3 weeks ago) bw_start_date_lp_3 = bw_start_date_lp_2 + relativedelta(weeks=-1) bw_end_date_lp_3 = bw_start_date_lp_3 + relativedelta(days=6) # LP (4 weeks ago) bw_start_date_lp_4 = bw_start_date_lp_3 + relativedelta(weeks=-1) bw_end_date_lp_4 = bw_start_date_lp_4 + relativedelta(days=6) print(f"Number of days in the building week = {n_building_days}") print(f"Today = {str(today)}") print(f"TP: {bw_start_date_tp} to {bw_end_date_tp}") print(f"LP_BW: {bw_start_date_lp_bw} to {bw_end_date_lp_bw}") print(f"LP_1: {bw_start_date_lp_1} to {bw_end_date_lp_1}") print(f"LP_2: {bw_start_date_lp_2} to {bw_end_date_lp_2}") print(f"LP_3: {bw_start_date_lp_3} to {bw_end_date_lp_3}") print(f"LP_4: {bw_start_date_lp_4} to {bw_end_date_lp_4}") #################################### print('Grabbing off limits labels...') off_limits_labels = snowflake_query('select * from awal.awal_ar.off_limits_labels') off_limits_labels['LABEL'] = '%' + off_limits_labels['LABEL'].str.lower() + '%' # create sql text for filtering out off limits labels i = 1 off_limits_labels_query = '' for each_name in off_limits_labels['LABEL']: off_limits_labels_query = off_limits_labels_query + f"lower(label) like $${each_name}$$ or lower(distributor) like $${each_name}$$" if i=1000*n_building_days) & (wtd['US_WTD_TP']>=1000*n_building_days)])==0: wtd = wtd.loc[(wtd['GL_WTD_TP']>=1000*n_building_days)].reset_index(drop=True) else: wtd = wtd.loc[(wtd['GL_WTD_TP']>=1000*n_building_days) & (wtd['US_WTD_TP']>=1000*n_building_days)].reset_index(drop=True) daily_streams = daily_streams.loc[daily_streams['UNIFIED_SONG_ID'].isin(wtd['UNIFIED_SONG_ID'])].reset_index(drop=True) # save here snowflake_query('delete from awal.awal_ar.DISCOVERY_DAILY_STREAMS') upload_df_to_snowflake_table(daily_streams,'DISCOVERY_DAILY_STREAMS') del daily_streams # getting building week growth (adjusting for denominator of 0; capping % at 999%) wtd['US_WTD_LP_ADJUSTED'] = wtd['US_WTD_LP'].clip(lower=0.0001) wtd['GL_WTD_LP_ADJUSTED'] = wtd['GL_WTD_LP'].clip(lower=0.0001) wtd['US_WTD_PCT'] = (wtd['US_WTD_TP'] - wtd['US_WTD_LP']) / wtd['US_WTD_LP_ADJUSTED'] wtd['US_WTD_PCT'] = wtd['US_WTD_PCT'].clip(upper=9.99) wtd['GL_WTD_PCT'] = (wtd['GL_WTD_TP'] - wtd['GL_WTD_LP']) / wtd['GL_WTD_LP_ADJUSTED'] wtd['GL_WTD_PCT'] = wtd['GL_WTD_PCT'].clip(upper=9.99) wtd.drop(['US_WTD_LP_ADJUSTED','GL_WTD_LP_ADJUSTED'],axis=1,inplace=True) ######## WEEKLY STREAMS ########## print('\nPull weekly streams...\n') weekly_streams = run_query(f""" SELECT unified_song_id, region, tw_on_demand_audio_streams weekly_streams, CASE WHEN report_date like $${bw_end_date_tp}$$ THEN 'LW1' WHEN report_date like $${bw_end_date_lp_1}$$ THEN 'LW2' WHEN report_date like $${bw_end_date_lp_2}$$ THEN 'LW3' WHEN report_date like $${bw_end_date_lp_3}$$ THEN 'LW4' END as WHICH_WEEK FROM luminate.vw_daily_top_songs_global WHERE 1=1 AND unified_song_id in {tuple(wtd['UNIFIED_SONG_ID'])} AND region IN ('US', 'Global') AND report_date in ($${bw_end_date_tp}$$,$${bw_end_date_lp_1}$$,$${bw_end_date_lp_2}$$,$${bw_end_date_lp_3}$$) and weekly_streams is not NULL ORDER BY unified_song_id, region,report_date ; """) weekly_streams.columns = weekly_streams.columns.str.upper() tmp_pivot = pd.pivot_table(weekly_streams, values='WEEKLY_STREAMS', index=['UNIFIED_SONG_ID','WHICH_WEEK'],columns=['REGION']).reset_index() tmp_pivot.fillna(0,inplace=True) del weekly_streams global_weekly_streams = organize_weekly_streams(tmp_pivot,'Global') us_weekly_streams = organize_weekly_streams(tmp_pivot,'US') # global_weekly_streams = tmp_pivot[['UNIFIED_SONG_ID','GL_LW1','GL_LW2','GL_LW3','GL_LW4']] # us_weekly_streams = tmp_pivot[['UNIFIED_SONG_ID','US_LW1','US_LW2','US_LW3','US_LW4']] weekly_streams = us_weekly_streams.merge(global_weekly_streams,how='left',on='UNIFIED_SONG_ID') weekly_streams['DATE_REPORTED'] = bw_end_date_tp weekly_streams['DATE_UPDATED'] = today # save here snowflake_query('delete from awal.awal_ar.DISCOVERY_WEEKLY_STREAMS') upload_df_to_snowflake_table(weekly_streams,'DISCOVERY_WEEKLY_STREAMS') del weekly_streams ################################### # get release dates ################################### alternate_isrcs = alternate_isrcs.loc[alternate_isrcs['UNIFIED_SONG_ID'].isin(wtd['UNIFIED_SONG_ID'])] ref_list = alternate_isrcs.loc[:][['ISRC']] still_need_metadata = ref_list.copy() # initial metadata pull print(f'Getting track metadata for {len(still_need_metadata)} ISRCs via spotify...') result_from_spotify = pull_metadata_from_spotify(still_need_metadata) # merge and ID which songs still missing ISRC metadata full_list = result_from_spotify.copy() still_need_metadata = ref_list.loc[~ref_list['ISRC'].isin(full_list['ISRC'])].reset_index(drop=True) print(f"Getting missing metadata for {len(still_need_metadata)} ISRCs via Apple/iTunes...") # get remaining metadata via apple/itunes if len(still_need_metadata)>0: result_from_itunes = pull_metadata_from_itunes(still_need_metadata) if len(result_from_itunes)>0: result_from_itunes = ref_list.merge(result_from_itunes,how='right',on='ISRC') # add fake columns to be able to concat result_from_itunes['SPOTIFY_TRACK_ID'] = None result_from_itunes['FOLLOWERS_LATEST'] = None else: result_from_itunes = pd.DataFrame() # add fake columns to be able to concat full_list['ITUNES_TRACK_ID'] = None full_list = pd.concat([full_list,result_from_itunes]).reset_index(drop=True) # ID which songs still missing metadata still_need_metadata = ref_list.loc[~ref_list['ISRC'].isin(full_list['ISRC'])].reset_index(drop=True) ################################### print(f"Getting missing metadata for {len(still_need_metadata)} ISRCs via Amazon...") # get remaining metadata via amazon if len(still_need_metadata)>0: result_from_amazon = pull_metadata_from_amazon(still_need_metadata) if len(result_from_amazon)>0: result_from_amazon = ref_list.merge(result_from_amazon,how='right',on='ISRC') # add fake columns to be able to concat result_from_amazon['SPOTIFY_TRACK_ID'] = None result_from_amazon['FOLLOWERS_LATEST'] = None result_from_amazon['ITUNES_TRACK_ID'] = None result_from_amazon['PREVIEW_URL'] = None else: result_from_amazon = pd.DataFrame() # add fake columns to be able to concat full_list['AMAZON_TRACK_ID'] = None full_list = pd.concat([full_list,result_from_amazon]).reset_index(drop=True) # ID which songs still missing metadata still_need_metadata = ref_list.loc[~ref_list['ISRC'].isin(full_list['ISRC'])].reset_index(drop=True) ################################### print(f"Getting missing metadata for {len(still_need_metadata)} ISRCs via Deezer...") # get remaining metadata via deezer if len(still_need_metadata)>0: result_from_deezer = pull_metadata_from_deezer(still_need_metadata) if len(result_from_deezer)>0: result_from_deezer = ref_list.merge(result_from_deezer,how='right',on='ISRC') # add fake columns to be able to concat result_from_deezer['SPOTIFY_TRACK_ID'] = None result_from_deezer['FOLLOWERS_LATEST'] = None result_from_deezer['ITUNES_TRACK_ID'] = None result_from_deezer['AMAZON_TRACK_ID'] = None else: result_from_deezer = pd.DataFrame() # add fake columns to be able to concat full_list['DEEZER_TRACK_ID'] = None full_list = pd.concat([full_list,result_from_deezer]).reset_index(drop=True) # ID which songs still missing metadata still_need_metadata = ref_list.loc[~ref_list['ISRC'].isin(full_list['ISRC'])].reset_index(drop=True) ################################### print(f"Getting missing metadata for {len(still_need_metadata)} ISRCs via New Luminate RDB...") if len(still_need_metadata)>0: # using try/except in case Luminate switches up this beta table try: my_query = f""" select r.isrc, s.recording_date release_date, s.genres, r.label_name + ', ' + s.label_name as label from luminate_raw.dim_recordings r left outer join luminate_raw.map_recording_songs ref on ref.mr_id=r.mr_id left outer join luminate_raw.dim_songs s on ref.song_id=s.song_id where 1=1 and r.isrc in {tuple(still_need_metadata['ISRC'])} """ result_from_rdb = run_query(my_query) result_from_rdb.columns = result_from_rdb.columns.str.upper() # extract genre result_from_rdb['GENRES'] = result_from_rdb['GENRES'].replace('None',None) result_from_rdb['GENRES'] = result_from_rdb['GENRES'].str.extract(r'MAIN_GENRE\"\:\"(.*?)\"', expand=False) if len(result_from_rdb)>0: result_from_rdb = ref_list.merge(result_from_rdb,how='right',on='ISRC') # add fake columns to be able to concat result_from_rdb['SPOTIFY_TRACK_ID'] = None result_from_rdb['FOLLOWERS_LATEST'] = None result_from_rdb['ITUNES_TRACK_ID'] = None result_from_rdb['AMAZON_TRACK_ID'] = None result_from_rdb['DEEZER_TRACK_ID'] = None result_from_rdb['PREVIEW_URL'] = None except: result_from_rdb = pd.DataFrame() else: result_from_rdb = pd.DataFrame() full_list = pd.concat([full_list,result_from_rdb]).reset_index(drop=True) # ID which songs still missing metadata still_need_metadata = ref_list.loc[~ref_list['ISRC'].isin(full_list['ISRC'])].reset_index(drop=True) ################################### print(f"Still missing metadata for {len(still_need_metadata)} ISRCS; no resolution here at the moment but will keep in final result") ################################### # generate track URL based on DSP full_list['TRACK_URL'] = np.where( ~full_list['SPOTIFY_TRACK_ID'].isna(), 'https://open.spotify.com/track/' + full_list['SPOTIFY_TRACK_ID'], np.where( ~full_list['ITUNES_TRACK_ID'].isna(), 'https://music.apple.com/us/song/' + full_list['ITUNES_TRACK_ID'], np.where( ~full_list['AMAZON_TRACK_ID'].isna(), 'https://music.amazon.com/tracks/' + full_list['AMAZON_TRACK_ID'], np.where( ~full_list['DEEZER_TRACK_ID'].isna(), 'https://www.deezer.com/us/track/' + full_list['DEEZER_TRACK_ID'].astype('str'), None ) ) ) ) full_list.drop(['SPOTIFY_TRACK_ID','ITUNES_TRACK_ID','AMAZON_TRACK_ID','DEEZER_TRACK_ID'],axis=1,inplace=True) ################################### # add in still_need_metadata if len(still_need_metadata)>0: full_list = still_need_metadata.merge(full_list,how='outer',on=['ISRC']) # grab UNIFIED_SONG_ID full_list = alternate_isrcs.merge(full_list,how='right',on='ISRC') # identify songs where an isrc is tagged with off limits labels again (each DSP sometimes has more info than luminate) off_limits_labels['LABEL'] = off_limits_labels['LABEL'].str.replace('%','') to_drop = full_list.loc[((full_list['LABEL'].str.lower()).fillna('').str.contains('|'.join(off_limits_labels['LABEL'])))]['UNIFIED_SONG_ID'] full_list.drop(full_list.loc[full_list['UNIFIED_SONG_ID'].isin(to_drop)].index,inplace=True) full_list.reset_index(drop=True,inplace=True) # keep the most relevant ISRC per unified song id full_list = full_list.sort_values(['UNIFIED_SONG_ID','RELEASE_DATE'],ignore_index=True) full_list.drop_duplicates('UNIFIED_SONG_ID',inplace=True) full_list.reset_index(drop=True,inplace=True) ################################### # merge to wtd print('Merging...') wtd.drop('ISRC',axis=1,inplace=True) final_list = wtd.merge(full_list,how='right',on='UNIFIED_SONG_ID') del wtd,full_list # fill in extra label info if new info final_list['LABEL_x'] = final_list['LABEL_x'].fillna('') final_list['LABEL_y'] = final_list['LABEL_y'].fillna('') final_list['LABEL'] = np.where( final_list.apply(lambda x: True if x['LABEL_y'] in x['LABEL_x'] else False,axis=1), final_list['LABEL_x'], final_list['LABEL_x'] + ', ' + final_list['LABEL_y'] ) final_list.drop(['LABEL_x','LABEL_y'],axis=1,inplace=True) final_list['GENRES'] = final_list['GENRES'].replace('{}',None) # fix Nan's vs. None vs. '' cols_to_fix = ['RELEASE_DATE','GENRES','TRACK_URL','LABEL'] for each_col in cols_to_fix: final_list[each_col] = final_list[each_col].astype('str').replace('nan',None) final_list[each_col] = final_list[each_col].astype('str').replace('None',None) final_list[each_col] = np.where( final_list[each_col].astype('str')=='', None, final_list[each_col] ) final_list.reset_index(drop=True,inplace=True) # drop irrelevant entries - we do this here otherwise they'll still exist in the result just with Nan's # Ensure RELEASE_DATE is datetime final_list['RELEASE_DATE'] = pd.to_datetime(final_list['RELEASE_DATE'], errors='coerce') print(len(final_list)) # drop a song if released prior to this century final_list = final_list.loc[(final_list['RELEASE_DATE'].isna()) | (final_list['RELEASE_DATE']>=pd.Timestamp('2000-01-01'))].reset_index(drop=True) print(len(final_list)) # contains these genres for genre_we_dont_want in off_limits_genres_contains: final_list = final_list.loc[(final_list['GENRES'].isna()) | (~final_list['GENRES'].fillna('').str.contains(genre_we_dont_want))].reset_index(drop=True) # equals this genre (single) for genre_we_dont_want in off_limits_genres_equals: final_list = final_list.loc[(final_list['GENRES'].isna()) | (final_list['GENRES'].fillna('')!=genre_we_dont_want)].reset_index(drop=True) # final_list = final_list.loc[(final_list['FOLLOWERS_LATEST'].isna()) | (final_list['FOLLOWERS_LATEST']<250000)].reset_index(drop=True) # drop artists (all songs) if artist's last release (or last relevant release after filtering) is prior to 2010 latest_metadata = final_list.copy() latest_metadata = latest_metadata.sort_values('RELEASE_DATE',ascending=False) latest_metadata = latest_metadata[[ 'UNIFIED_ARTIST_ID', 'RELEASE_DATE' ]] latest_metadata.drop_duplicates('UNIFIED_ARTIST_ID',inplace=True) latest_metadata = latest_metadata[ latest_metadata['RELEASE_DATE'].isna() | (latest_metadata['RELEASE_DATE'] >= pd.Timestamp('2010-01-01')) ] final_list = final_list.loc[final_list['UNIFIED_ARTIST_ID'].isin(latest_metadata['UNIFIED_ARTIST_ID'])] del latest_metadata print(len(final_list)) final_list['RELEASE_DATE'] = final_list['RELEASE_DATE'].dt.date # convert back to date final_list.drop(['FOLLOWERS_LATEST'],axis=1,inplace=True) final_list['DATE_REPORTED'] = bw_end_date_tp final_list['DATE_UPDATED'] = today print(f"Final result: {len(final_list.loc[final_list['TRACK_URL'].isna() & final_list['RELEASE_DATE'].isna()])} songs without metadata") ################################### # save here to a main table snowflake_query('delete from awal.awal_ar.DISCOVERY_MAIN') upload_df_to_snowflake_table(final_list,'DISCOVERY_MAIN') print(f"Finished. Time elapsed = {(dt.datetime.now() - start_time).seconds/60:.2f} minutes // {len(final_list)} songs uploaded")