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 time
import snowflake.connector
from snowflake.connector.pandas_tools import write_pandas
import os
import boto3
import redshift_connector
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']
# except:
# pass
try:
env = os.getenv("Environment", "dev")
succeeded = True
except:
env = "dev"
succeeded = False
session = boto3.session.Session()
client = session.client(
service_name='secretsmanager',
region_name='us-east-1'
)
sf_user = client.get_secret_value(SecretId=f"{env}/awal-ar/SNOWFLAKE_USER")["SecretString"]
sf_account = client.get_secret_value(SecretId=f"{env}/awal-ar/SNOWFLAKE_ACCOUNT")["SecretString"]
sf_warehouse = client.get_secret_value(SecretId=f"{env}/awal-ar/SNOWFLAKE_WAREHOUSE")["SecretString"]
sf_password = client.get_secret_value(SecretId=f"{env}/awal-ar/PEM_KEY_PASSWORD")["SecretString"]
sf_pem_key = client.get_secret_value(SecretId=f"{env}/awal-ar/PEM_KEY")["SecretString"]
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 setup_logger():
import logging
logger = logging.getLogger()
logger.setLevel(logging.INFO)
# Add a logger handler if none exists
if not logger.hasHandlers():
handler = logging.StreamHandler()
formatter = logging.Formatter('%(asctime)s - %(levelname)s - %(message)s')
handler.setFormatter(formatter)
logger.addHandler(handler)
return logger
def orcd_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 delphi_query(sql_query):
ctx = snowflake.connector.connect(
user=sf_user,
private_key=pkb,
account='delphi',
warehouse='AWAL_ANALYTICS_LARGE_WAREHOUSE',
region='us-east-1'
)
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 db_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
# temporary stage so we can just join it to the next queries instead of using "where x in {}"
def upload_temp_stage_to_delphi(stage_name,stage_data):
delphi_query(f'REMOVE @~/{stage_name}.csv.gz')
with open(f'{stage_name}.csv', 'w') as f:
for each_item in stage_data:
f.write(f"{each_item}\n")
file_path = os.path.abspath(f'{stage_name}.csv')
delphi_query(f'PUT file://{file_path} @~ AUTO_COMPRESS=TRUE')
def filter_spotify():
# 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
# track level data only
track_result = orcd_query(f"""
select
distinct
mp.song_id,
t.SPOTIFY_TRACK_ID,
coalesce(a.copyrights[1]:"text"::string,'') as LABEL,
a.release_date,
t.PREVIEW_URL,
t.duration_ms
from chartmetric.raw_data.spotify t
join chartmetric.raw_data.spotify_album a on a.SPOTIFY_ALBUM_ID=t.spotify_album_id
join awal.awal_ar.MAP_SONG_ID_SPOTIFY_TRACK_ID mp on mp.SPOTIFY_TRACK_ID = t.SPOTIFY_TRACK_ID
where 1=1
--and t.preview_url is not null
""")
# duration
off_limits_dur = track_result.loc[(track_result['DURATION_MS']<60000) | (track_result['DURATION_MS']>420000)]['SONG_ID']
track_result.drop('DURATION_MS',axis=1,inplace=True)
# release date - ID dates that are before 2000 and also are already out of bounds (like year 1000). Let Nan's stay
# track_result['RELEASE_DATE'] = pd.to_datetime(track_result['RELEASE_DATE'])
# off_limits_release_date = track_result.loc[track_result['RELEASE_DATE']
SPLIT(t.SPOTIFY_ARTIST_ID, ',')) AS split_spotify_artist_id
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 TRIM(split_spotify_artist_id.value::string) = artist.SPOTIFY_ARTIST_ID
join awal.awal_ar.MAP_SONG_ID_SPOTIFY_TRACK_ID mp on mp.SPOTIFY_TRACK_ID = t.SPOTIFY_TRACK_ID
""")
artist_result.drop_duplicates(inplace=True)
# contains these collabs
off_limits_coll = artist_result.loc[artist_result['SPOTIFY_ARTIST_NAME'].str.lower().str.contains('|'.join(off_limits_collabs['ARTIST']),na=False)]['SONG_ID']
# equals these artists
off_limits_art_0 = artist_result.loc[artist_result['SPOTIFY_ARTIST_NAME'].str.lower().isin(off_limits_artists['ARTIST'].tolist())]['SONG_ID']
off_limits_art_1 = artist_result.loc[artist_result['SPOTIFY_ARTIST_NAME'].isin(['Joey Bada$$','¥$'])]['SONG_ID']
# contains these genres
off_limits_gen_0 = artist_result.loc[artist_result['GENRES'].str.lower().str.contains('|'.join(off_limits_genres_contains['GENRE']),na=False)]['SONG_ID']
# equals this genre (single)
genres_we_dont_want = off_limits_genres_equals['GENRE'].tolist()
off_limits_gen_1 = artist_result.loc[artist_result['GENRES'].str.lower().isin(genres_we_dont_want)]['SONG_ID']
off_limits = pd.concat([
off_limits_dur,
off_limits_release_date,
off_limits_lab_0,
off_limits_lab_1,
off_limits_lab_2,
off_limits_lab_3,
off_limits_coll,
off_limits_art_0,
off_limits_art_1,
off_limits_gen_0,
off_limits_gen_1
],ignore_index=True)
off_limits.drop_duplicates(inplace=True)
if len(off_limits)>0:
track_result = track_result.drop(track_result.loc[track_result['SONG_ID'].isin(off_limits)].index).reset_index(drop=True)
artist_result = artist_result.drop(artist_result.loc[artist_result['SONG_ID'].isin(off_limits)].index).reset_index(drop=True)
return off_limits,track_result,artist_result
def remove_irrelevant_titles(result):
# drop if title has White Noise and genre is NEARS. And so on
keywords_to_exclude = [
'Noise','Rain','Sleep','Relax','Hz','Waves','Ocean','Calm','Sounds','Slumber','Peace','Tranquil'
]
result = result[
~(
result['TITLE'].str.contains('|'.join(keywords_to_exclude)) & (result['GENRE']=='NEARS')
)
].reset_index(drop=True)
result = result[
~(
result['LABEL'].str.contains('Platoon ',na=False) & result['GENRE'].str.contains('NEARS',na=False)
)
].reset_index(drop=True)
# just question marks or spaces
result = result[~result['TITLE'].str.contains(r'^[? ]*$')].reset_index(drop=True)
# metadata = metadata[~metadata['SONG_TITLE'].str.contains(r'(?=.*Hz)(?=.*Frequency)')].reset_index(drop=True)
keywords_to_exclude = [
' Hz',
'Rain Sound',
'Spa Music',
'Twinkle Twinkle Little Star',
'Itsy Bitsy Spider',
'Hush Little Baby',
'Brown Noise'
]
result = result[~result['TITLE'].str.contains('|'.join(keywords_to_exclude))].reset_index(drop=True)
# removing entries with special characters
for each_col in ['ARTIST','TITLE']:
# just a mix of question marks + other symbols
result = result[~result[each_col].str.match(r'[-+?¿.&,()\\~ ]*$')].reset_index(drop=True)
# song title has quotation mark
result = result[~result['TITLE'].str.match(r'.*”.*')].reset_index(drop=True)
# artist/title is just spaces
result = result.loc[~result[each_col].str.match(r'[(\u200e)]*$')].reset_index(drop=True)
# artist/title STARTS with a random symbol
result = result[~result[each_col].str.match(r'^[,`├┘─┬┼╤╨╫]')].reset_index(drop=True)
# starts with [?]
result = result[~result[each_col].str.match(r'^\[\?\]')].reset_index(drop=True)
# artist has 3+ question marks
result = result[~result['ARTIST'].str.match(r'.*\?.*\?.*\?.*')].reset_index(drop=True)
# song has 4+ question marks
result = result[~result['TITLE'].str.match(r'.*\?.*\?.*\?.*\?')].reset_index(drop=True)
# at least 3 upside down question marks
result = result[~result[each_col].str.match(r'.*\¿.*\¿.*\¿.*')].reset_index(drop=True)
# arabic, hebrew
# result = result[~result[each_col].str.match(r'.*[\u0600-\u06FF].*')].reset_index(drop=True)
# result = result[~result[each_col].str.match(r'.*[\u0590-\u05FF].*')].reset_index(drop=True)
# thai
result = result[~result[each_col].str.match(r'.*[\u0E00-\u0E7F].*')].reset_index(drop=True)
# other non-linguistic characters
result = result[~result[each_col].str.match(r'.*[🎵�ª¦‡§▓▒░¬╝»『』┼╜╨╣╤╗╕╬├║╡《》].*')].reset_index(drop=True)
# foreign languages
# result = result[~result[each_col].str.match(r'.*[țĂșłğŽΤΦΩΓλдшрНЦбЛйάαлДพื้นทีที่ซ้อะฟกบวว่าភក្ដីស្នបណ្តូលពេจคำเชยๆமறைகிறாய்नमालिङ्गाष्टकम्मीशमीशाननिर्वाणरूपंรัญជ្រេហ៍ချို့ჩუბδინააჭარულიმელოდიებიโ].*')].reset_index(drop=True)
# korean
# result = result[~result[each_col].str.match(r'.*[\u1100-\u11FF\u3130-\u318F\uAC00-\uD7AF].*')].reset_index(drop=True)
# chinese/japanese
# result = result[~result[each_col].str.match(r'.*[\u4e00-\u9fff].*')].reset_index(drop=True)
# result = result[~result[each_col].str.match(r'.*[\u3040-\u309F\u30A0-\u30FF\u4E00-\u9FFF].*')].reset_index(drop=True)
# vietnamese / portuguese
# result = result[~result[each_col].str.match(r'.*[ĩứưẤÅåắăặạậệâỗốờơộșȘ±¨].*')].reset_index(drop=True)
# starts with a dash or plus or Δ, followed by special character
result = result[~result[each_col].str.match(r'^[-+Δ][¥Ò£îº¡\?ç]')].reset_index(drop=True)
# # Regex pattern to match names entirely in a language
# for pattern in [
# # r'^[\u4e00-\u9fff]+$', # chinese
# # r'^[\u1100-\u11FF\u3130-\u318F\uAC00-\uD7AF]+$', # korean
# # r'^[\u3040-\u309F\u30A0-\u30FF\u4E00-\u9FFF]+$' # japanese
# # r'^[\u0E00-\u0E7F]+$', # thai
# r'^[\u0370-\u03FF]+$', # greek
# r'^[\u0400-\u04FF]+$', # russian
# # r'^[\u0590-\u05FF]+$', # hebrew
# # r'^[\u0600-\u06FF]+$', # arabic
# r'^[\u0100-\u024F]+$', # Romanian / turkish
# r'^[\u10A0-\u10FF]+$', # Georgian
# r'^[\u0530-\u058F]+$', # armenian
# r'^[\u1780-\u17FF]+$', # cambodian
# r'^[\u1000-\u109F]+$', # burmese
# r'^[\u0B80-\u0BFF]+$', # tamil
# r'^[\u0900-\u097F]+$' # Hindi, sanskrit
# ]:
# result = result.loc[~result['ARTIST'].str.match(pattern)]
# result = result.loc[~result['TITLE'].str.match(pattern)]
return result
# drop artists (all songs) if artist's last release (or last relevant release after filtering) is prior to 2015
def release_date_filtering(df,which_col):
# drop artists (all songs) if artist's last release (or last relevant release after filtering) is prior to 2015
latest_metadata = df.copy()
latest_metadata = latest_metadata.sort_values('RELEASE_DATE',ascending=False)
latest_metadata = latest_metadata[[
which_col,
'RELEASE_DATE'
]]
latest_metadata.drop_duplicates(which_col,inplace=True)
to_drop = latest_metadata[
(~latest_metadata['RELEASE_DATE'].isna()) &
(latest_metadata['RELEASE_DATE'] < pd.Timestamp('2015-01-01'))
]
to_drop = to_drop[which_col]
return to_drop
def combine_labels(df,colA,colB):
# Fill nulls and strip whitespace
df[colA] = df[colA].fillna('').str.strip()
df[colB] = df[colB].fillna('').str.strip()
# Lowercase versions for case-insensitive check
label_x_lower = df[colA].str.lower()
label_y_lower = df[colB].str.lower()
# Check substrings row-wise
x_in_y = [x in y for x, y in zip(label_x_lower, label_y_lower)]
y_in_x = [y in x for x, y in zip(label_x_lower, label_y_lower)]
# Convert to Series
x_in_y = pd.Series(x_in_y, index=df.index)
y_in_x = pd.Series(y_in_x, index=df.index)
# Use np.select as before
df['LABEL'] = np.select(
[
x_in_y & (df[colB].str.len() > df[colA].str.len()),
y_in_x & (df[colA].str.len() >= df[colB].str.len())
],
[
df[colB],
df[colA]
],
default=df[colA] + ', ' + df[colB]
)
return df
def filter_redundant_substrings(distributor_list):
# Clean up strings: remove duplicates, trim whitespace
cleaned = list(set(s.strip() for s in distributor_list if pd.notna(s)))
result = []
for d in cleaned:
# Only add if it's not a substring of something already in result
if not any(d.lower() in other.lower() and d != other for other in cleaned):
result.append(d)
return ', '.join(sorted(result)) # sort optional, for consistency
logger = setup_logger()
logger.info('Starting main script...')
if succeeded:
logger.info(f"ENV = {env}")
# set start time just so we can see how long this takes to run
start_time = dt.datetime.now()
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)
logger.info(f"Number of days in the building week = {n_building_days}")
logger.info(f"Today = {str(today)}")
logger.info(f"TP: {bw_start_date_tp} to {bw_end_date_tp}")
logger.info(f"LP_BW: {bw_start_date_lp_bw} to {bw_end_date_lp_bw}")
logger.info(f"LP_1: {bw_start_date_lp_1} to {bw_end_date_lp_1}")
logger.info(f"LP_2: {bw_start_date_lp_2} to {bw_end_date_lp_2}")
logger.info(f"LP_3: {bw_start_date_lp_3} to {bw_end_date_lp_3}")
logger.info(f"LP_4: {bw_start_date_lp_4} to {bw_end_date_lp_4}")
####################################
logger.info('Grabbing off limits labels...')
off_limits_labels = orcd_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) not like $$%{each_name}%$$"
if i= '2000-01-01'))
and ((label is NULL) or ({off_limits_labels_query}))
and ((label is NULL) or (
upper(label) not in (
$$P2025$$,
$$P2024$$,
$$P2023$$,
$$P2022$$,
$$P2021$$,
$$P2020$$,
$$WAR$$,
$$COLU$$,
$$RHI$$,
$$DEF$$,
$$WARN$$,
$$UNIV$$,
$$INT$$,
$$UME$$,
$$COL$$,
$$RHIN$$,
$$DISN$$,
$$POLY$$,
$$APG1$$,
$$GEFN$$,
$$ARCTIC$$,
$$AVIE RECORDS$$,
$$MIRARE$$,
$$SPAR$$,
$$SOUNDSCAPE$$,
$$ISLAND$$,
$$RHINO$$,
$$GIANT$$
)
and label not like $$RCA%$$
))
),
artist_ids as (
SELECT song_id, f.value:ARTIST_ID::STRING AS artist_id
FROM LUMINATE_DB_LISTING_DETAIL.EXTRACT_S.VW_SONG_DS s
LEFT JOIN LATERAL FLATTEN(INPUT => PARSE_JSON(COALESCE(s.ARTISTS, '[]'::VARCHAR))) f
),
wtd as (
SELECT
d.song_id,
SUM(CASE WHEN d.country_code = 'AA' THEN d.quantity ELSE 0 END) AS gl_wtd,
SUM(CASE WHEN d.country_code = 'US' THEN d.quantity ELSE 0 END) AS us_wtd,
max(d.report_date) as report_date
FROM LUMINATE_DB_LISTING_DETAIL.EXTRACT_S.VW_DAILY_FACT_SONG_SUMMARY_DS d
where 1=1
and d.country_code in ('AA','US')
and d.metric_category like 'Streams'
and d.content_type like 'Audio'
and d.service_type like 'OnDemand'
and d.report_date between $${bw_start_date_tp}$$ and $${bw_end_date_tp}$$
group by 1
),
prefinal as (
select
wtd.*,
artist_ids.* exclude(song_id),
song_ids.* exclude(song_id)
from wtd
join
song_ids on song_ids.song_id=wtd.song_id
LEFT JOIN -- get artist id using json; have to do this here otherwise songs w/o artist id will not appear
artist_ids
ON artist_ids.song_id = wtd.song_id
left join song_ids_from_watch_list s on wtd.song_id=s.song_id
where (
--us_wtd > {n_building_days} * 500
gl_wtd > {n_building_days} * 1000
) or (
s.song_id IS NOT NULL
)
),
artist_ids_who_havent_released_since as (
SELECT distinct ARTIST_ID
from prefinal
group by artist_id
having max(release_date) < '2015-01-01'
),
artist_names_who_havent_released_since as (
SELECT distinct lower(ARTIST) as artist
from prefinal
where ARTIST_ID is NULL
group by lower(artist)
having max(release_date) < '2015-01-01'
)
select *
from prefinal
where 1=1
and (
artist_id is not null and NOT EXISTS (
SELECT 1
FROM artist_ids_who_havent_released_since a1
WHERE prefinal.ARTIST_ID = a1.artist_id
)
or (
artist_id is null and NOT EXISTS (
SELECT 1
FROM artist_names_who_havent_released_since a2
WHERE lower(prefinal.artist) = a2.artist
)
)
)
""")
logger.info(f"Initial pull = {len(df)}")
df = df.loc[~df['SONG_ID'].isin(olsi['SONG_ID'])]
logger.info(f"Removed off-limits song id's = {len(df)}")
logger.info(f"Time elapsed = {(dt.datetime.now() - start_time).seconds/60:.2f} minutes")
logger.info('Filtering...')
df = remove_irrelevant_titles(df)
date_reported = df['REPORT_DATE'].max()
# get isrcs of initial result so we can use this to filter out off limits distros
logger.info('Getting isrcs...')
upload_temp_stage_to_delphi('song_ids',df['SONG_ID'].drop_duplicates())
isrcs = delphi_query(f"""
with song_ids as (
SELECT $1 AS song_id FROM @~/song_ids.csv.gz
)
select distinct song_ids.song_id,mr.isrc
from song_ids
left join LUMINATE_DB_LISTING_DETAIL.EXTRACT_S.VW_MR_SONG_MAP_DS ref on ref.song_id=song_ids.song_id
left join LUMINATE_DB_LISTING_DETAIL.EXTRACT_S.VW_MUSICAL_RECORDING_DS mr on ref.mr_id=mr.mr_id
where mr.isrc is not null and mr.type like 'Audio'
""")
upload_temp_stage_to_delphi('isrcs',isrcs['ISRC'].drop_duplicates())
orcd_query('delete from awal.awal_ar.ISRC_MAP_SONG_ID')
upload_df_to_snowflake_table(isrcs,'ISRC_MAP_SONG_ID')
del isrcs
logger.info('Pull distro\'s...')
try:
test_pull = delphi_query(f"""
select "label" from LUMINATE_MODELS.PROD.DIM_GFK_METADATA_RDB limit 1
""")
pull_distro = True
except:
pull_distro = False
if pull_distro:
logger.info("Pull distro metadata directly from Delphi...")
distributor_metadata = delphi_query(f"""
with isrcs as (
SELECT $1 AS isrc FROM @~/isrcs.csv.gz
),
max_date as (
select max(last_isrc_date) latest_date from LUMINATE_MODELS.PROD.DIM_GFK_METADATA_RDB
)
select
g.isrc,
CASE
WHEN g."label" IS NULL THEN g.distributor
WHEN g.distributor IS NULL THEN g."label"
WHEN POSITION(lower(g."label") IN lower(g.distributor)) > 0 THEN g.distributor
WHEN POSITION(lower(g.distributor) IN lower(g."label")) > 0 THEN g."label"
ELSE g."label" || ', ' || g.distributor -- If no overlap
END AS distributor
from isrcs i
join LUMINATE_MODELS.PROD.DIM_GFK_METADATA_RDB g on g.isrc=i.isrc
join max_date on max_date.latest_date=g.LAST_ISRC_DATE
where g.product_type like 'audio'
""")
distributor_metadata = distributor_metadata.loc[distributor_metadata['DISTRIBUTOR']!="\Apple Music DJ Mixes"]
distributor_metadata = distributor_metadata.loc[distributor_metadata['DISTRIBUTOR']!="Apple Music"]
orcd_query("delete from awal.awal_ar.distributors")
upload_df_to_snowflake_table(distributor_metadata[['ISRC','DISTRIBUTOR']],'DISTRIBUTORS')
del distributor_metadata
else:
logger.info("Using most recent distro metadata...")
# remove off limits distro's
logger.info("Filtering distro's...")
distributors = orcd_query(f"""
select
distinct
mp.song_id,
d.isrc,
d.distributor
from awal.awal_ar.distributors d
join awal.awal_ar.ISRC_MAP_SONG_ID mp on mp.isrc=d.isrc
""")
off_limits_distributors = distributors.loc[distributors['DISTRIBUTOR'].str.lower().str.contains('|'.join(off_limits_labels['LABEL']))]
if len(off_limits_distributors)>0:
df = df.drop(df.loc[df['SONG_ID'].isin(off_limits_distributors['SONG_ID'].drop_duplicates())].index).reset_index(drop=True)
distributors = distributors.drop(distributors.loc[distributors['SONG_ID'].isin(off_limits_distributors['SONG_ID'].drop_duplicates())].index).reset_index(drop=True)
distributors.drop('ISRC',axis=1,inplace=True)
distributors.drop_duplicates(inplace=True)
# aggregate
distributors = (
distributors
.drop_duplicates(subset=['SONG_ID', 'DISTRIBUTOR']) # optional, safe
.groupby('SONG_ID')['DISTRIBUTOR']
.apply(filter_redundant_substrings)
.reset_index()
)
# filter off limits via spotify metadata
logger.info('Getting spotify track IDs...')
upload_temp_stage_to_delphi('song_ids',df['SONG_ID'].drop_duplicates())
# get all spotify track IDs per song id
spotify_track_ids_map = delphi_query(f"""
with song_ids as (
SELECT $1 AS song_id
FROM @~/song_ids.csv.gz
)
select
song_ids.song_id,
f.value::string AS SPOTIFY_TRACK_ID
from LUMINATE_DB_LISTING_DETAIL.EXTRACT_S.VW_SONG_DS s
join LATERAL FLATTEN(input => s.EXTERNAL_IDS:SPOTIFY) f
join song_ids on song_ids.song_id=s.song_id
""")
orcd_query('delete from awal.awal_ar.MAP_SONG_ID_SPOTIFY_TRACK_ID')
upload_df_to_snowflake_table(spotify_track_ids_map,'MAP_SONG_ID_SPOTIFY_TRACK_ID')
del spotify_track_ids_map
logger.info(f"Time elapsed = {(dt.datetime.now() - start_time).seconds/60:.2f} minutes")
logger.info('Filtering spotify...')
off_limits_spotify,metadata_spotify_track,metadata_spotify_artist = filter_spotify()
if len(off_limits_spotify)>0:
df = df.drop(df.loc[df['SONG_ID'].isin(off_limits_spotify)].index).reset_index(drop=True)
del off_limits_spotify
logger.info(f"Time elapsed = {(dt.datetime.now() - start_time).seconds/60:.2f} minutes")
logger.info('Filtering via release date...')
# drop artists (all songs) if artist's last release (or last relevant release after filtering) is prior to 2015
# Ensure RELEASE_DATE is datetime
df['RELEASE_DATE'] = pd.to_datetime(df['RELEASE_DATE'], errors='coerce')
metadata_spotify_track['RELEASE_DATE'] = pd.to_datetime(metadata_spotify_track['RELEASE_DATE'], errors='coerce')
# ARTIST_ID (release date via luminate)
to_drop = release_date_filtering(df.loc[~df['ARTIST_ID'].isna()],'ARTIST_ID')
df = df[~df['ARTIST_ID'].isin(to_drop)]
# SPOTIFY_ARTIST_ID (release date via spotify)
tmp = metadata_spotify_track[['SONG_ID','RELEASE_DATE']].merge(metadata_spotify_artist[['SONG_ID','SPOTIFY_ARTIST_ID']],how='left',on='SONG_ID')
to_drop = release_date_filtering(tmp.loc[~tmp['SPOTIFY_ARTIST_ID'].isna()],'SPOTIFY_ARTIST_ID')
to_drop_track_id = metadata_spotify_artist[metadata_spotify_artist['SPOTIFY_ARTIST_ID'].isin(to_drop)]['SONG_ID']
metadata_spotify_artist = metadata_spotify_artist[~metadata_spotify_artist['SPOTIFY_ARTIST_ID'].isin(to_drop)]
metadata_spotify_track = metadata_spotify_track[~metadata_spotify_track['SONG_ID'].isin(to_drop_track_id)]
df = df[~df['SONG_ID'].isin(to_drop_track_id)]
# ARTIST (name)
to_drop = release_date_filtering(df.loc[~df['ARTIST'].isna()],'ARTIST')
df = df[~df['ARTIST'].isin(to_drop)]
logger.info('Merging metadata...')
# merge spotify metadata
df = df.merge(metadata_spotify_track[['SONG_ID','SPOTIFY_TRACK_ID','RELEASE_DATE','LABEL','PREVIEW_URL']],how='left',on='SONG_ID')
df = df.merge(metadata_spotify_artist[['SONG_ID','SPOTIFY_ARTIST_ID','SPOTIFY_ARTIST_NAME','GENRES','ARTIST_ORDER']],how='left',on='SONG_ID')
del metadata_spotify_track,metadata_spotify_artist
# fill in extra info if new info
## SPOTIFY
df = combine_labels(df,'LABEL_x','LABEL_y')
df.drop(['LABEL_x','LABEL_y'],axis=1,inplace=True)
# genres
df['GENRES'] = df['GENRES'].replace('{}',None)
df['GENRE'] = df['GENRE'].fillna('')
df['GENRES'] = df['GENRES'].fillna('')
df['GENRE'] = np.where(
df.apply(lambda x: True if x['GENRE'] in x['GENRES'] else False,axis=1),
df['GENRES'],
df['GENRE'] + ', ' + df['GENRES']
)
df.drop(['GENRES'],axis=1,inplace=True)
# release dates
df['RELEASE_DATE'] = np.where(
df['RELEASE_DATE_x'].isna(),
df['RELEASE_DATE_y'],
df['RELEASE_DATE_x']
)
df.drop(['RELEASE_DATE_x','RELEASE_DATE_y'],axis=1,inplace=True)
## Merge DISTRO metadata
df = df.merge(distributors,how='left',on='SONG_ID')
# label
df = combine_labels(df,'LABEL','DISTRIBUTOR')
df.drop(['DISTRIBUTOR'],axis=1,inplace=True)
df.replace('',None,inplace=True)
df['LABEL'] = df['LABEL'].str.replace('n.a., ', '')
df['LABEL'] = df['LABEL'].str.replace(', n.a.', '')
df['LABEL'] = df['LABEL'].str.replace('** UNKNOWN **, ', '')
df['LABEL'] = df['LABEL'].str.replace(', ** UNKNOWN **', '')
# remove off limits labels again now that combined
off_limits = df.loc[df['LABEL'].str.lower().str.contains('|'.join(off_limits_labels['LABEL']),na=False)]
if len(off_limits)>0:
df = df.loc[~df['SONG_ID'].isin(off_limits['SONG_ID'])].reset_index(drop=True)
logger.info(f"Time elapsed = {(dt.datetime.now() - start_time).seconds/60:.2f} minutes")
logger.info('Pulling Daily streams...')
# get wtd LP
song_ids = df.drop_duplicates('SONG_ID')['SONG_ID'].tolist()
upload_temp_stage_to_delphi('song_ids',song_ids)
daily_streams = delphi_query(f"""
with song_ids as (
SELECT $1 AS song_id FROM @~/song_ids.csv.gz
),
pre as (
SELECT
d.song_id,
d.report_date,
sum(CASE WHEN d.country_code = 'AA' THEN d.quantity ELSE 0 END) AS gl_daily,
sum(CASE WHEN d.country_code = 'US' THEN d.quantity ELSE 0 END) AS us_daily,
CASE
WHEN d.report_date like $${bw_end_date_tp - relativedelta(days=6)}$$ THEN 'DAY_1'
WHEN d.report_date like $${bw_end_date_tp - relativedelta(days=5)}$$ THEN 'DAY_2'
WHEN d.report_date like $${bw_end_date_tp - relativedelta(days=4)}$$ THEN 'DAY_3'
WHEN d.report_date like $${bw_end_date_tp - relativedelta(days=3)}$$ THEN 'DAY_4'
WHEN d.report_date like $${bw_end_date_tp - relativedelta(days=2)}$$ THEN 'DAY_5'
WHEN d.report_date like $${bw_end_date_tp - relativedelta(days=1)}$$ THEN 'DAY_6'
WHEN d.report_date like $${bw_end_date_tp}$$ THEN 'DAY_7'
ELSE NULL
END as WHICH_DAY
from LUMINATE_DB_LISTING_DETAIL.EXTRACT_S.VW_DAILY_FACT_SONG_SUMMARY_DS d
join song_ids tmp on tmp.song_id=d.song_id
where 1=1
and d.country_code in ('AA','US')
and d.metric_category like 'Streams'
and d.content_type like 'Audio'
and d.service_type like 'OnDemand'
and d.report_date between $${bw_start_date_lp_bw}$$ and $${bw_end_date_tp}$$
group by 1,2
order by d.song_id,d.report_date
),
wtd_lp as (
select
song_id,
sum(GL_DAILY) as GL_WTD_LP,
sum(US_DAILY) as US_WTD_LP
from pre
where 1=1
and report_date <= $${bw_end_date_lp_bw}$$
group by song_id
),
adjusted as (
select
song_id,
GL_DAILY,
US_DAILY,
WHICH_DAY
from pre
where 1=1
and which_day is not null
),
pivoted as (
SELECT
song_id,
COALESCE(SUM(CASE WHEN which_day = 'DAY_1' THEN us_daily ELSE 0 END), 0) AS us_day_1,
COALESCE(SUM(CASE WHEN which_day = 'DAY_2' THEN us_daily ELSE 0 END), 0) AS us_day_2,
COALESCE(SUM(CASE WHEN which_day = 'DAY_3' THEN us_daily ELSE 0 END), 0) AS us_day_3,
COALESCE(SUM(CASE WHEN which_day = 'DAY_4' THEN us_daily ELSE 0 END), 0) AS us_day_4,
COALESCE(SUM(CASE WHEN which_day = 'DAY_5' THEN us_daily ELSE 0 END), 0) AS us_day_5,
COALESCE(SUM(CASE WHEN which_day = 'DAY_6' THEN us_daily ELSE 0 END), 0) AS us_day_6,
COALESCE(SUM(CASE WHEN which_day = 'DAY_7' THEN us_daily ELSE 0 END), 0) AS us_day_7,
COALESCE(SUM(CASE WHEN which_day = 'DAY_1' THEN gl_daily ELSE 0 END), 0) AS gl_day_1,
COALESCE(SUM(CASE WHEN which_day = 'DAY_2' THEN gl_daily ELSE 0 END), 0) AS gl_day_2,
COALESCE(SUM(CASE WHEN which_day = 'DAY_3' THEN gl_daily ELSE 0 END), 0) AS gl_day_3,
COALESCE(SUM(CASE WHEN which_day = 'DAY_4' THEN gl_daily ELSE 0 END), 0) AS gl_day_4,
COALESCE(SUM(CASE WHEN which_day = 'DAY_5' THEN gl_daily ELSE 0 END), 0) AS gl_day_5,
COALESCE(SUM(CASE WHEN which_day = 'DAY_6' THEN gl_daily ELSE 0 END), 0) AS gl_day_6,
COALESCE(SUM(CASE WHEN which_day = 'DAY_7' THEN gl_daily ELSE 0 END), 0) AS gl_day_7
FROM adjusted
GROUP BY song_id
ORDER BY song_id
)
select
pivoted.*,
wtd_lp.US_WTD_LP,
wtd_lp.GL_WTD_LP
from pivoted
left join wtd_lp on wtd_lp.song_id=pivoted.song_id
""")
wtd_lp = daily_streams[['SONG_ID','GL_WTD_LP','US_WTD_LP']]
daily_streams.drop(['GL_WTD_LP','US_WTD_LP'],axis=1,inplace=True)
daily_streams['DATE_REPORTED'] = date_reported
daily_streams['DATE_UPDATED'] = today
# save here
orcd_query('delete from awal.awal_ar.DISCOVERY_DAILY_STREAMS')
upload_df_to_snowflake_table(daily_streams,'DISCOVERY_DAILY_STREAMS')
del daily_streams
## merge wtd_lp into df
df = df.merge(wtd_lp,how='left',on='SONG_ID')
df[['GL_WTD_LP','US_WTD_LP']] = df[['GL_WTD_LP','US_WTD_LP']].fillna(0)
df.rename(columns={'GL_WTD':'GL_WTD_TP','US_WTD':'US_WTD_TP'},inplace=True)
df['GL_WTD_CHG'] = df['GL_WTD_TP'] - df['GL_WTD_LP']
df['US_WTD_CHG'] = df['US_WTD_TP'] - df['US_WTD_LP']
# getting building week growth (adjusting for denominator of 0; capping % at 999%)
df['US_WTD_LP_ADJUSTED'] = df['US_WTD_LP'].clip(lower=0.0001)
df['GL_WTD_LP_ADJUSTED'] = df['GL_WTD_LP'].clip(lower=0.0001)
df['US_WTD_PCT'] = (df['US_WTD_TP'] - df['US_WTD_LP']) / df['US_WTD_LP_ADJUSTED']
df['US_WTD_PCT'] = df['US_WTD_PCT'].clip(upper=9.99)
df['GL_WTD_PCT'] = (df['GL_WTD_TP'] - df['GL_WTD_LP']) / df['GL_WTD_LP_ADJUSTED']
df['GL_WTD_PCT'] = df['GL_WTD_PCT'].clip(upper=9.99)
df.drop(['US_WTD_LP_ADJUSTED','GL_WTD_LP_ADJUSTED'],axis=1,inplace=True)
logger.info(f"Time elapsed = {(dt.datetime.now() - start_time).seconds/60:.2f} minutes")
######## WEEKLY STREAMS ##########
logger.info('\nPull weekly streams...\n')
weekly_streams = delphi_query(f"""
with song_ids as (
SELECT $1 AS song_id FROM @~/song_ids.csv.gz
)
SELECT
d.song_id,
country_code region,
dt.WEEKID,
sum(quantity) as streams
from LUMINATE_DB_LISTING_DETAIL.EXTRACT_S.VW_DAILY_FACT_SONG_SUMMARY_DS d
join LUMINATE_DB_LISTING_DETAIL.EXTRACT_S.VW_DATE_DS dt on d.report_date = TO_DATE(TO_VARCHAR(dt.dateid), 'YYYYMMDD')
join song_ids tmp on tmp.song_id=d.song_id
where 1=1
and d.country_code in ('AA','US')
and d.metric_category like 'Streams'
and d.content_type like 'Audio'
and d.service_type like 'OnDemand'
and d.report_date between $${bw_start_date_lp_4}$$ and $${bw_end_date_lp_1}$$
group by 1,2,3
order by 1,2,3
""")
weekids = weekly_streams['WEEKID'].drop_duplicates().sort_values().reset_index(drop=True)
weekly_streams['WHICH_WEEK'] = np.where(
weekly_streams['WEEKID']==weekids[0],
'LW4',
np.where(
weekly_streams['WEEKID']==weekids[1],
'LW3',
np.where(
weekly_streams['WEEKID']==weekids[2],
'LW2',
'LW1'
)
)
)
weekly_streams.drop('WEEKID',axis=1,inplace=True)
weekly_streams['REGION'] = weekly_streams['REGION'].replace('AA','GL')
# Pivot the DataFrame
pivoted = weekly_streams.pivot(index='SONG_ID', columns=['REGION','WHICH_WEEK'], values=['STREAMS'])
# Flatten MultiIndex columns
pivoted.columns = [f"{col[1]}_{col[2]}" for col in pivoted.columns]
# Fill missing with 0
pivoted = pivoted.fillna(0)
# Reset index to make song_id a column again
pivoted = pivoted.reset_index()
for col in pivoted.columns[1:]:
pivoted[col] = pivoted[col].astype('int64')
weekly_streams = pivoted.copy()
weekly_streams['DATE_REPORTED'] = date_reported
weekly_streams['DATE_UPDATED'] = today
# save here
orcd_query('delete from awal.awal_ar.DISCOVERY_WEEKLY_STREAMS')
upload_df_to_snowflake_table(weekly_streams,'DISCOVERY_WEEKLY_STREAMS')
del weekly_streams
logger.info(f"Time elapsed = {(dt.datetime.now() - start_time).seconds/60:.2f} minutes")
logger.info('Last bit of organizing...')
df['TRACK_URL'] = np.where(
~df['SPOTIFY_TRACK_ID'].isna(),
'https://open.spotify.com/track/' + df['SPOTIFY_TRACK_ID'],
None
)
df.rename(columns={'REPORT_DATE':'DATE_REPORTED'},inplace=True)
df['RELEASE_DATE'] = df['RELEASE_DATE'].dt.date
df['DATE_UPDATED'] = today
df = df[[
'SONG_ID','ARTIST_ID','SPOTIFY_ARTIST_ID',
'ARTIST','TITLE','SPOTIFY_ARTIST_NAME','ARTIST_ORDER',
'US_WTD_TP','US_WTD_LP','GL_WTD_TP','GL_WTD_LP',
'US_WTD_CHG','US_WTD_PCT','GL_WTD_CHG','GL_WTD_PCT',
'RELEASE_DATE','GENRE','TRACK_URL','PREVIEW_URL','LABEL',
'DATE_REPORTED','DATE_UPDATED'
]]
## Fix missing data on artist_id or spotify_artist_id
# Step 1: Build mapping from known pairs
id_pairs = df.dropna(subset=['ARTIST_ID', 'SPOTIFY_ARTIST_ID']).drop_duplicates()
artist_to_spotify = dict(zip(id_pairs['ARTIST_ID'], id_pairs['SPOTIFY_ARTIST_ID']))
spotify_to_artist = dict(zip(id_pairs['SPOTIFY_ARTIST_ID'], id_pairs['ARTIST_ID']))
# Step 2: Fill missing SPOTIFY_ARTIST_ID only if artist_order == 0
df['SPOTIFY_ARTIST_ID'] = df.apply(
lambda row: artist_to_spotify.get(row['ARTIST_ID'], row['SPOTIFY_ARTIST_ID'])
if row.get('artist_order') == 0 else row['SPOTIFY_ARTIST_ID'],
axis=1
)
# Step 3: Fill missing ARTIST_ID only if artist_order == 0
df['ARTIST_ID'] = df.apply(
lambda row: spotify_to_artist.get(row['SPOTIFY_ARTIST_ID'], row['ARTIST_ID'])
if row.get('artist_order') == 0 else row['ARTIST_ID'],
axis=1
)
logger.info(f"Final result: {len(df.drop_duplicates('SONG_ID'))} songs; {len(df)} total rows")
###################################
# save here to a main table
orcd_query('delete from awal.awal_ar.DISCOVERY_MAIN')
upload_df_to_snowflake_table(df,'DISCOVERY_MAIN')
logger.info(f"Time elapsed = {(dt.datetime.now() - start_time).seconds/60:.2f} minutes")
###################################
# remove artists if their latest release is off limits.
orcd_query(f"""
EXECUTE TASK dev_engineering.jho.remove_artists_by_latest_release;
""")
time.sleep(10)
###################################
# artist table - this will also trigger subsequent tasks
orcd_query(f"""
EXECUTE TASK dev_engineering.jho.discovery_artist_table;
""")
###################################
logger.info('Getting artist metadata...')
artist_ids = df['ARTIST_ID'].dropna().drop_duplicates()
artist_data = delphi_query(f"""
select
artist_id,
artist_name as unified_artist_name,
artist_type,
music_type,
COUNTRY_OF_ORIGIN
from LUMINATE_DB_LISTING_DETAIL.EXTRACT_S.VW_ARTIST_DS
where artist_id in {tuple(artist_ids)}
""")
# save here
orcd_query('delete from awal.awal_ar.discovery_artist_metadata')
upload_df_to_snowflake_table(artist_data,'DISCOVERY_ARTIST_METADATA')
logger.info(f"Finished. Time elapsed = {(dt.datetime.now() - start_time).seconds/60:.2f} minutes")