from dateutil.relativedelta import relativedelta, FR import datetime as dt import pytz import requests import streamlit as st import pandas as pd import numpy as np import snowflake.connector from snowflake.connector.pandas_tools import write_pandas import os import boto3 def get_access_token(login_username,login_password,api_key): url = "https://api.luminatedata.com/auth" payload = { "username": login_username, "password": login_password } headers = { "accept": "application/json", "x-api-key": api_key, "content-type": "application/x-www-form-urlencoded" } response = requests.post(url, data=payload, headers=headers) resp_json = response.json() acc_token = resp_json["access_token"] return acc_token def get_song_streams_json(isrc,region,start_date,end_date,login_username,api_key,acc_token,frequency): if frequency=='DAILY': url = f"https://api.luminatedata.com/songs/{isrc}?id_type=isrc&location={region}&start_date={start_date}&end_date={end_date}&external_ids=true&relationships=song&metrics=streams&breakouts=all&aggregate_interval=day&markets=&market_filter=&commercial_model=none&service_type=on_demand&content_type=audio" elif frequency=='WEEKLY': new_start_date = start_date + relativedelta(weeks=-4,weekday=FR) # 5 weeks ago new_end_date = start_date + relativedelta(days=6) url = f"https://api.luminatedata.com/songs/{isrc}?id_type=isrc&location={region}&start_date={new_start_date}&end_date={new_end_date}&external_ids=true&relationships=song&metrics=streams&breakouts=all&aggregate_interval=chart_week&markets=&market_filter=&commercial_model=none&service_type=on_demand&content_type=audio" headers = { "email": login_username, "Accept": "application/vnd.luminate-data.svc-apibff.v1+json", "x-api-key": api_key, "authorization": acc_token } response = requests.get(url, headers=headers) try: resp_json = response.json() except: resp_json = None return resp_json def organize_song_streams_data(isrc,resp_json,region,wtd_start_date,frequency): # get artist id try: artist_id = resp_json['artists'][0]['id'] except: artist_id = None # get "on-demand" results if exists try: streams_response = resp_json['metrics'][0]['value'][3]['value'][0] if streams_response['name']=='on_demand': streams_response = streams_response['value'] else: streams_response = None except: streams_response = None # organize results if frequency=='DAILY': if streams_response is not None: lw_daily_dates = [item['date'] for item in streams_response if item['date']=str(wtd_start_date)] tw_daily_streams = [item['value'] for item in streams_response if item['date']>=str(wtd_start_date)] n_building_days = len(tw_daily_streams) lw_final = sum(lw_daily_streams) wtd_lp = sum(lw_daily_streams[0:n_building_days]) wtd_tp = sum(tw_daily_streams) # merge all the daily numbers together daily_dates = lw_daily_dates + tw_daily_dates daily_streams = lw_daily_streams + tw_daily_streams else: daily_dates = None daily_streams = None lw_final = None wtd_lp = None wtd_tp = None if daily_dates is not None and len(daily_dates) > 1: item_result = pd.DataFrame({ 'ISRC':isrc, 'REGION':region, 'ARTIST_ID':artist_id, 'DAILY_DATES':daily_dates, 'DAILY_STREAMS':daily_streams, 'LW_FINAL':lw_final, 'WTD_LP':wtd_lp, 'WTD_TP':wtd_tp }) else: item_result = pd.DataFrame({ 'ISRC':[isrc], 'REGION':[region], 'ARTIST_ID':[artist_id], 'DAILY_DATES':daily_dates, 'DAILY_STREAMS':daily_streams, 'LW_FINAL':[lw_final], 'WTD_LP':[wtd_lp], 'WTD_TP':[wtd_tp] }) elif frequency=='WEEKLY': if streams_response is not None: weekly_dates = [item['end_date'] for item in streams_response] weekly_streams = [item['value'] for item in streams_response] else: weekly_dates = None weekly_streams = None if weekly_dates is not None and len(weekly_dates) > 1: item_result = pd.DataFrame({ 'ISRC':isrc, 'REGION':region, 'WEEKLY_DATES':weekly_dates, 'WEEKLY_STREAMS':weekly_streams, }) else: item_result = pd.DataFrame({ 'ISRC':[isrc], 'REGION':[region], 'WEEKLY_DATES':weekly_dates, 'WEEKLY_STREAMS':weekly_streams, }) return item_result def get_song_atd(isrc,region,end_date,login_username,api_key,acc_token): if region=='US': start_date = '2013-12-30' elif region=='AA': start_date = '2019-01-04' url = f"https://api.luminatedata.com/songs/{isrc}?id_type=isrc&location={region}&start_date={start_date}&end_date={end_date}&metrics=streams&breakouts=all&aggregate_interval=total&service_type=on_demand&content_type=audio" headers = { "email": login_username, "Accept": "application/vnd.luminate-data.svc-apibff.v1+json", "x-api-key": api_key, "authorization": acc_token } response = requests.get(url, headers=headers) # organize # try: resp_json = response.json() streams_response = resp_json['metrics'][0]['value'][3]['value'][0] if streams_response['name']=='on_demand': streams_response = streams_response['value'] else: streams_response = None except: streams_response = None if streams_response is not None: atd = streams_response else: atd = None item_result = pd.DataFrame({ 'ISRC':isrc, 'REGION':region, 'ATD':[atd] }) return item_result def get_artist_streams_json(artist_id,region,start_date,end_date,login_username,api_key,acc_token,frequency): if frequency=='DAILY': url = f"https://api.luminatedata.com/artists/{artist_id}?id_type=luminate&location={region}&start_date={start_date}&end_date={end_date}&relationships=artist&metrics=streams&breakouts=all&aggregate_interval=day&service_type=on_demand&content_type=audio&sales_type=&store_strata=" elif frequency=='WEEKLY': new_start_date = start_date + relativedelta(weeks=-4,weekday=FR) # 5 weeks ago new_end_date = start_date + relativedelta(days=6) url = f"https://api.luminatedata.com/artists/{artist_id}?id_type=luminate&location={region}&start_date={new_start_date}&end_date={new_end_date}&relationships=artist&metrics=streams&breakouts=all&aggregate_interval=chart_week&service_type=on_demand&content_type=audio&sales_type=&store_strata=" headers = { "email": login_username, "Accept": "application/vnd.luminate-data.svc-apibff.v1+json", "x-api-key": api_key, "authorization": acc_token } response = requests.get(url, headers=headers) try: resp_json = response.json() except: resp_json = None return resp_json def organize_artist_streams_data(artist_id,resp_json,region,wtd_start_date,frequency): # get "on-demand" results if exists try: streams_response = resp_json['metrics'][0]['value'][3]['value'][0] if streams_response['name']=='on_demand': streams_response = streams_response['value'] else: streams_response = None except: streams_response = None # organize results if frequency=='DAILY': if streams_response is not None: lw_daily_dates = [item['date'] for item in streams_response if item['date']=str(wtd_start_date)] tw_daily_streams = [item['value'] for item in streams_response if item['date']>=str(wtd_start_date)] n_building_days = len(tw_daily_streams) lw_final = sum(lw_daily_streams) wtd_lp = sum(lw_daily_streams[0:n_building_days]) wtd_tp = sum(tw_daily_streams) # merge all the daily numbers together daily_dates = lw_daily_dates + tw_daily_dates daily_streams = lw_daily_streams + tw_daily_streams else: daily_dates = None daily_streams = None lw_final = None wtd_lp = None wtd_tp = None if daily_dates is not None and len(daily_dates) > 1: item_result = pd.DataFrame({ 'REGION':region, 'ARTIST_ID':artist_id, 'DAILY_DATES':daily_dates, 'DAILY_STREAMS':daily_streams, 'LW_FINAL':lw_final, 'WTD_LP':wtd_lp, 'WTD_TP':wtd_tp }) else: item_result = pd.DataFrame({ 'REGION':[region], 'ARTIST_ID':[artist_id], 'DAILY_DATES':daily_dates, 'DAILY_STREAMS':daily_streams, 'LW_FINAL':[lw_final], 'WTD_LP':[wtd_lp], 'WTD_TP':[wtd_tp] }) elif frequency=='WEEKLY': if streams_response is not None: weekly_dates = [item['end_date'] for item in streams_response] weekly_streams = [item['value'] for item in streams_response] else: weekly_dates = None weekly_streams = None if weekly_dates is not None and len(weekly_dates) > 1: item_result = pd.DataFrame({ 'REGION':region, 'ARTIST_ID':artist_id, 'WEEKLY_DATES':weekly_dates, 'WEEKLY_STREAMS':weekly_streams, }) else: item_result = pd.DataFrame({ 'REGION':[region], 'ARTIST_ID':[artist_id], 'WEEKLY_DATES':weekly_dates, 'WEEKLY_STREAMS':weekly_streams, }) return item_result def get_artist_atd(artist_id,region,end_date,login_username,api_key,acc_token): if region=='US': start_date = '2013-12-30' elif region=='AA': start_date = '2019-01-04' url = f"https://api.luminatedata.com/artists/{artist_id}?id_type=luminate&location={region}&start_date={start_date}&end_date={end_date}&relationships=artist&metrics=streams&breakouts=all&aggregate_interval=total&service_type=on_demand&content_type=audio&sales_type=&store_strata=" headers = { "email": login_username, "Accept": "application/vnd.luminate-data.svc-apibff.v1+json", "x-api-key": api_key, "authorization": acc_token } response = requests.get(url, headers=headers) # organize # try: resp_json = response.json() streams_response = resp_json['metrics'][0]['value'][3]['value'][0] if streams_response['name']=='on_demand': streams_response = streams_response['value'] else: streams_response = None except: streams_response = None if streams_response is not None: atd = streams_response else: atd = None item_result = pd.DataFrame({ 'ARTIST_ID':artist_id, 'REGION':region, 'ATD':[atd] }) return item_result 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'] login_password = os.environ['LUMINATE_PASSWORD'] api_key = os.environ['LUMINATE_API_KEY'] 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'] login_password = st.secrets['LUMINATE_PASSWORD'] api_key = st.secrets['LUMINATE_API_KEY'] 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"] login_password = client.get_secret_value(SecretId="dev/awal-ar/LUMINATE_PASSWORD")["SecretString"] api_key = client.get_secret_value(SecretId="dev/awal-ar/LUMINATE_API_KEY")["SecretString"] except: pass login_username = sf_user 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 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() def get_isrcs(): df = snowflake_query('select * from awal.awal_ar.song_list')[['ISRC']] isrcs = df.loc[(~df['ISRC'].isna()) & (df['ISRC']!='')] isrcs = isrcs.drop_duplicates('ISRC')['ISRC'].values.tolist() return isrcs today = dt.datetime.now(pytz.timezone('US/Eastern')).date() if today.weekday() == 5: wtd_start_date = today + relativedelta(days=-1, weekday=FR(-2)) else: wtd_start_date = today + relativedelta(days=-1, weekday=FR(-1)) start_date = wtd_start_date - relativedelta(days=7) # fri lw end_date = today + relativedelta(days=-2) # latest reported date isrcs = get_isrcs() df = [] # authenticate luminate / get access token acc_token = get_access_token(login_username,login_password,api_key) print('Access token = ' + acc_token) print(f'{len(isrcs)} isrcs to pull') print('Pulling song streams...') song_streams_data = [] counter_i = 1 for isrc in isrcs: print(f"ISRC {counter_i} of {len(isrcs)}") for region in ['US','AA']: for freq in ['DAILY','WEEKLY']: resp_json = get_song_streams_json(isrc,region,start_date,end_date,login_username,api_key,acc_token,freq) if resp_json is not None: item_result = organize_song_streams_data(isrc,resp_json,region,wtd_start_date,freq) song_streams_data.append(item_result) counter_i = counter_i+1 song_streams_df = pd.concat(song_streams_data).reset_index(drop=True) print('Pulling song ATD...') song_atd_data = [] counter_i = 1 for isrc in isrcs: print(f"ISRC {counter_i} of {len(isrcs)}") for region in ['US','AA']: item_result = get_song_atd(isrc,region,end_date,login_username,api_key,acc_token) song_atd_data.append(item_result) counter_i = counter_i+1 song_atd_df = pd.concat(song_atd_data).reset_index(drop=True) artist_ids = song_streams_df.loc[~song_streams_df['ARTIST_ID'].isna()]['ARTIST_ID'].drop_duplicates() print('Pulling artist streams...') artist_streams_data = [] counter_i = 1 for artist_id in artist_ids: print(f"Artist {counter_i} of {len(artist_ids)}") for region in ['US','AA']: for freq in ['DAILY','WEEKLY']: resp_json = get_artist_streams_json(artist_id,region,start_date,end_date,login_username,api_key,acc_token,freq) if resp_json is not None: item_result = organize_artist_streams_data(artist_id,resp_json,region,wtd_start_date,freq) artist_streams_data.append(item_result) counter_i = counter_i+1 artist_streams_df = pd.concat(artist_streams_data).reset_index(drop=True) print('Pulling artist ATD...') artist_atd_data = [] counter_i = 1 for artist_id in artist_ids: print(f"Artist {counter_i} of {len(artist_ids)}") for region in ['US','AA']: item_result = get_artist_atd(artist_id,region,end_date,login_username,api_key,acc_token) artist_atd_data.append(item_result) counter_i = counter_i+1 artist_atd_df = pd.concat(artist_atd_data).reset_index(drop=True) # merge daily streams with atd data song_df = song_streams_df.merge(song_atd_df,how='left',on=['ISRC','REGION']) artist_df = artist_streams_df.merge(artist_atd_df,how='left',on=['ARTIST_ID','REGION']) # dates song_df['DATE_REPORTED'] = song_df['DAILY_DATES'].dropna().max() song_df['DATE_UPDATED'] = today artist_df['DATE_REPORTED'] = artist_df['DAILY_DATES'].dropna().max() artist_df['DATE_UPDATED'] = today # replace AA with GL song_df['REGION'] = np.where( song_df['REGION']=='AA', 'GL', song_df['REGION'] ) artist_df['REGION'] = np.where( artist_df['REGION']=='AA', 'GL', artist_df['REGION'] ) # save to snowflake snowflake_query('delete from awal.awal_ar.SONG_STREAMS') upload_df_to_snowflake_table(song_df,'SONG_STREAMS') snowflake_query('delete from awal.awal_ar.ARTIST_STREAMS') upload_df_to_snowflake_table(artist_df,'ARTIST_STREAMS') print('Finished')