# This file contains all the logics generate the data for agg_user_daily table. from datetime import datetime, timedelta, timezone import pandas as pd import snowflake.snowpark as snowpark agg_play_info_hourly = [ "CLIP_ID", "USER_ID", "USER_UID", "PLAY_DURATION_SEC", "PLAY_CNT", "PLAY_DURATION_THRESHOLD", "IS_USER_SONG_OWNER", "PLATFORM", "P_DATE", "P_HOUR", ] def get_sql_result(session: snowpark.Session, schema: str): df = session.sql(schema) local_df = df.collect() if len(local_df) == 0: raise Exception("No data found") return pd.DataFrame(local_df) def create_temp_table(session: snowpark.Session, cur_p_date: str, cur_p_hour: int, pre_p_date: str, pre_p_hour: int): session.sql(""" CREATE OR REPLACE TRANSIENT TABLE temp_agg_play_info AS select '{p_date}' as p_date, {p_hour} as p_hour, '{pre_p_date}' as pre_p_date, {pre_p_hour} as pre_p_hour; """.format(p_date=cur_p_date, p_hour=cur_p_hour, pre_p_date=pre_p_date, pre_p_hour=pre_p_hour)).collect() # process the two_hour_ago data only def get_agg_web_play_info( session: snowpark.Session, cur_p_date: str, cur_p_hour: int, play_duration_threshold: int = 0 ): schema = """ with a as ( select song_session_id, song_id as clip_id, play_duration, user_id, is_user_song_owner, action_index, client_timestamp, from web_audio_player_actions where play_duration >= 0 and user_id is not null and song_session_id != '' and DATE(p_date) = '{p_date}' and p_hour = {p_hour} ) select a.clip_id, e.user_id as user_id, a.user_id as user_uid, sum(play_duration) as play_duration_sec, count(distinct song_session_id) as play_cnt, {play_duration_threshold} as play_duration_threshold, a.is_user_song_owner as is_user_song_owner, 'web' as platform, '{p_date}' as p_date, {p_hour} as p_hour from a left join rds_discord_info as e on a.user_id = e.uid where play_duration > {play_duration_threshold} group by a.user_id, a.clip_id, e.user_id, a.is_user_song_owner; """.format(p_date=cur_p_date, p_hour=cur_p_hour, play_duration_threshold=play_duration_threshold) print(schema) return get_sql_result(session, schema) # process the two_hour_ago android data only def get_agg_android_play_info( session: snowpark.Session, cur_p_date: str, cur_p_hour: int, play_duration_threshold: int = 0 ): schema = """ with a as ( select user_id, element_id as clip_id, parse_json(context) as context from app_event where os_name = 'Android' and CATEGORY = 'audio_player' and context is not null and user_id is not null and length(user_id) > 0 and p_date = '{p_date}' and p_hour = {p_hour} ), c as ( with b as ( select a.clip_id, CASE WHEN LOWER(a.context:isUserSongOwner) = 'true' THEN 2 WHEN LOWER(a.context:isUserSongOwner) = 'false' THEN 1 WHEN a.context:isUserSongOwner IS NULL THEN 1 END AS is_user_song_owner from a ) select clip_id, CASE when MAX(is_user_song_owner) = 2 then true else false end as is_user_song_owner from b group by clip_id ) select a.clip_id, b.user_id as user_id, a.user_id as user_uid, sum(a.context:playDuration) as play_duration_sec, count(distinct a.context:songSessionId) as play_cnt, {play_duration_threshold} as play_duration_threshold, c.is_user_song_owner as is_user_song_owner, 'android' as platform, '{p_date}' as p_date, {p_hour} as p_hour from a left join rds_discord_info as b on a.user_id = b.uid left join c on a.clip_id = c.clip_id where a.context:playDuration > {play_duration_threshold} and a.context:playDuration is not null and a.context:playDuration != 'null' and a.context:playDuration < 100000 group by b.user_id, a.clip_id, a.user_id, is_user_song_owner; """.format(p_date=cur_p_date, p_hour=cur_p_hour, play_duration_threshold=play_duration_threshold) print(schema) return get_sql_result(session, schema) # process the two_hour_ago ios data only def get_agg_ios_play_info_v2( session: snowpark.Session, cur_p_date: str, cur_p_hour: int, play_duration_threshold: int = 0 ): schema = """ with a as ( select user_id, element_id as clip_id, parse_json(context) as context from app_event where os_name = 'iOS' and CATEGORY = 'omni_player' and context is not null and user_id is not null and length(user_id) > 0 and p_date = '{p_date}' and p_hour = {p_hour} ), c as ( with b as ( select a.clip_id, CASE WHEN LOWER(a.context:isUserSongOwner) = 'true' THEN 2 WHEN LOWER(a.context:isUserSongOwner) = 'false' THEN 1 WHEN a.context:isUserSongOwner IS NULL THEN 1 END AS is_user_song_owner from a ) select clip_id, CASE when MAX(is_user_song_owner) = 2 then true else false end as is_user_song_owner from b group by clip_id ) select a.clip_id, b.user_id as user_id, a.user_id as user_uid, sum(a.context:playDuration) as play_duration_sec, count(distinct a.context:songSessionId) as play_cnt, {play_duration_threshold} as play_duration_threshold, c.is_user_song_owner as is_user_song_owner, 'ios' as platform, '{p_date}' as p_date, {p_hour} as p_hour from a left join rds_discord_info as b on a.user_id = b.uid left join c on a.clip_id = c.clip_id where a.context:playDuration > {play_duration_threshold} and a.context:playDuration is not null and a.context:playDuration != 'null' and a.context:playDuration < 100000 group by b.user_id, a.clip_id, a.user_id, is_user_song_owner; """.format(p_date=cur_p_date, p_hour=cur_p_hour, play_duration_threshold=play_duration_threshold) print(schema) return get_sql_result(session, schema) def merge_datafromes(session: snowpark.Session, p_date: str, p_hour: int): input = datetime.strptime(p_date, "%Y-%m-%d") + timedelta(hours=p_hour) one_hour_ago_date = p_date one_hour_ago_hour = p_hour create_temp_table(session, p_date, p_hour, one_hour_ago_date, one_hour_ago_hour) df1 = get_agg_web_play_info(session, one_hour_ago_date, one_hour_ago_hour, play_duration_threshold=0) df2 = get_agg_web_play_info(session, one_hour_ago_date, one_hour_ago_hour, play_duration_threshold=5) df5 = get_agg_android_play_info( session, one_hour_ago_date, one_hour_ago_hour, play_duration_threshold=0 ) df6 = get_agg_android_play_info( session, one_hour_ago_date, one_hour_ago_hour, play_duration_threshold=5 ) df7 = get_agg_ios_play_info_v2( session, one_hour_ago_date, one_hour_ago_hour, play_duration_threshold=0 ) df8 = get_agg_ios_play_info_v2( session, one_hour_ago_date, one_hour_ago_hour, play_duration_threshold=5 ) df = pd.concat([df1, df2, df5, df6, df7, df8], ignore_index=True) session.create_dataframe(df).write.save_as_table("agg_play_info_hourly_v0", mode="append") if validate_date(df): print("Data validation passed.") return slack_alert_query = """ select POST_TO_SLACK('#tech-alerts', 'agg_play_info_hourly_v0_insert: failed validation') """ session.sql(slack_alert_query).collect() def validate_date(df): # Check for null values in USER_UID null_user_uid_cnt = df["USER_UID"].isnull().sum() print (f"Null USER_UID count: {null_user_uid_cnt}") if null_user_uid_cnt > 0: raise Exception("there are null user_uid in the data") # Check for duplicates based on specified columns # duplicates = df.duplicated(subset=["CLIP_ID", "USER_ID", "USER_UID", "PLAY_DURATION_THRESHOLD", "PLATFORM", "IS_USER_SONG_OWNER"], keep=False) # non_unique_rows = df[duplicates] # print (f"Non-unique rows count: {len(non_unique_rows)}") # if len(non_unique_rows) > 0: # raise Exception("there are duplicates in the data") # Check user_id null rates is lower than 1% null_user_id_cnt = df["USER_UID"].isnull().sum() print (f"Null USER_ID count: {null_user_id_cnt}") if null_user_id_cnt / len(df) > 0.1: raise Exception("there are more than 1% null user_id in the data") # Check clip_id null rates is lower than 1% null_clip_id_cnt = df["CLIP_ID"].isnull().sum() print (f"Null CLIP_ID count: {null_clip_id_cnt}") if null_clip_id_cnt / len(df) > 0.1: raise Exception("there are more than 1% null clip_id in the data") # Check play_duration_sec is greater than 3600 second # greater_than_3600 = df[df["PLAY_DURATION_SEC"] < 3600] # print (f"Rows with PLAY_DURATION_SEC < 3600: {len(greater_than_3600)}") # if len(greater_than_3600) > 0: # return False return True def backfill_data_agg_play_info_hourly(session: snowpark.Session): start = datetime.strptime("2025-07-06 12:00:00", "%Y-%m-%d %H:%M:%S") end = datetime.strptime("2025-07-06 13:00:00", "%Y-%m-%d %H:%M:%S") current = start while current < end: for p_hour in range(0, 1): now = current + timedelta(hours=p_hour) if now > end: break print(now) one_hour_ago_date = (now - timedelta(hours=1)).strftime("%Y-%m-%d") one_hour_ago_hour = (now - timedelta(hours=1)).hour df1 = get_agg_web_play_info( session, one_hour_ago_date, one_hour_ago_hour, play_duration_threshold=0 ) df2 = get_agg_web_play_info( session, one_hour_ago_date, one_hour_ago_hour, play_duration_threshold=5 ) df5 = get_agg_android_play_info( session, one_hour_ago_date, one_hour_ago_hour, play_duration_threshold=0 ) df6 = get_agg_android_play_info( session, one_hour_ago_date, one_hour_ago_hour, play_duration_threshold=5 ) df7 = get_agg_ios_play_info_v2( session, one_hour_ago_date, one_hour_ago_hour, play_duration_threshold=0 ) df8 = get_agg_ios_play_info_v2( session, one_hour_ago_date, one_hour_ago_hour, play_duration_threshold=5 ) df = pd.concat([df1, df2, df5, df6, df7, df8], ignore_index=True) session.create_dataframe(df).write.save_as_table( "agg_play_info_hourly_v0", mode="append" ) print ("finish " + str(now)) current = current + timedelta(days=1)