# This file contains all the logics generate the data for agg_user_daily table. import snowflake.snowpark as snowpark TABLE_NAME = "FACT_PLAY" TASK_NAME = "FACT_PLAY_hourly_upsert" THRESHOLD_PERCENTAGE = 0.1 # Custom threshold for play duration checks PLAY_DURATION_THRESHOLD_PERCENTAGE = 0.8 # More strict threshold for play duration PLAY_DURATION_MIN_COUNT_THRESHOLD = 300000 MONITOR_TABLE_COLUMNS = [ "table_name", "task_name", "p_date", "p_hour", "validation_name", "validation_value", "validation_result", ] TODAY_COUNT_COLUMN_NAME = "TODAY_COUNT" YESTERDAY_COUNT_COLUMN_NAME = "YESTERDAY_COUNT" QUALITY_CHECK_THRESHOLD = 0 UUID_LENGTH = 36 LONG_UUID_LENGTH = 64 HOURLY_OCCURENCE_COLUMN_NAME = "HOURLY_OCCURENCE" TOTAL_COUNT_COLUMN_NAME = "TOTAL_COUNT" MIN_COUNT_THRESHOLD = 100000 # Default threshold volume checks DATA_VOLUME_CHECKS = { "total_row_count": None, "distinct_user_uid": "USER_UID", "distinct_user_id": "USER_ID", } # Custom threshold volume checks for play duration PLAY_DURATION_VOLUME_CHECKS = { "sum_play_duration": "PLAY_DURATION_SECONDS", } SINGLE_FIELD_QUALITY_CHECKS = { "song_session_id": { "column_names": ["SONG_SESSION_ID"], "checks": [ { "name": "null", "threshold": 0.02, }, { "name": "length", "length": UUID_LENGTH, "threshold": 0.001, }, ], }, "user_uid": { "column_names": ["USER_UID"], "checks": [ { "name": "length", "length": UUID_LENGTH, "threshold": 0.001, }, { "name": "null", "threshold": 0.001, }, ], }, "play_start_ts_utc": { "column_names": ["PLAY_START_TS_UTC"], "checks": [ { "name": "null", "threshold": 0.001, }, ], }, "play_end_ts_utc": { "column_names": ["PLAY_END_TS_UTC"], "checks": [ { "name": "null", "threshold": 0.001, }, ], }, "clip_duration_seconds": { "column_names": ["CLIP_DURATION_SECONDS"], "checks": [ { "name": "negative", "threshold": 0.001, }, ], }, "play_duration_seconds": { "column_names": ["PLAY_DURATION_SECONDS"], "checks": [ { "name": "negative", "threshold": 0.001, }, { "name": "zero", "threshold": 0.001, }, ], }, "country": { "column_names": ["COUNTRY"], "checks": [ { "name": "null", "threshold": 0.05, }, ], }, "session_id": { "column_names": ["SESSION_ID"], "checks": [ { "name": "length", "length": LONG_UUID_LENGTH, "threshold": 0.001, }, ], } } GROUP_BY_COLUMNS = [ "PLATFORM", "PLAY_PAGE_TYPE", "IS_MOBILE", ] def main(session: snowpark.Session, p_date: str, p_hour: int): import importlib data_validation = importlib.import_module("data_validation") validate_data = data_validation.validate_data quality_checks = SINGLE_FIELD_QUALITY_CHECKS # Validate data with default threshold validate_data( session, p_date, p_hour, TABLE_NAME, TASK_NAME, MONITOR_TABLE_COLUMNS, DATA_VOLUME_CHECKS, quality_checks, GROUP_BY_COLUMNS, THRESHOLD_PERCENTAGE, MIN_COUNT_THRESHOLD, ) # Validate play duration data with custom threshold validate_data( session, p_date, p_hour, TABLE_NAME, TASK_NAME, MONITOR_TABLE_COLUMNS, PLAY_DURATION_VOLUME_CHECKS, {}, # No quality checks needed for this run GROUP_BY_COLUMNS, PLAY_DURATION_THRESHOLD_PERCENTAGE, PLAY_DURATION_MIN_COUNT_THRESHOLD, )