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")