import warnings from datetime import datetime import dagster as dg from dagster import EnvVar from dagster_snowflake import SnowflakeResource # Using asset key references to avoid import chain issues from src.utils.snowflake.constants import Group, PartitionExpr, Warehouse from src.utils.snowflake.query import JinjaSQLFormatter, PythonStringSQLFormatter from src.utils.automation_conditions import daily_cron_with_eager_historical_backfill_condition from src.utils.glue_utils import trigger_glue_job warnings.filterwarnings("ignore", category=dg.BetaWarning) USER_TOP_SCORE_SONGS_START_DATE = datetime.strptime("2025-08-19", "%Y-%m-%d") # Last 30 days USER_TOP_SCORE_SONGS_TABLE_NAME = "USER_TOP_SCORE_SONGS" class UserTopScoreSongsConfig(dg.Config): user_top_score_songs_table_name: str = USER_TOP_SCORE_SONGS_TABLE_NAME warehouse: str = Warehouse.SUNO_PROD_USER_TOP_SCORE_SONGS.value @dg.asset( name="user_top_score_songs_snowflake", description="Daily user top scoring songs based on play count, play duration, and like count with ranking.", group_name=Group.PERSONALIZATION.value, partitions_def=dg.DailyPartitionsDefinition( start_date=USER_TOP_SCORE_SONGS_START_DATE, end_offset=0 ), # non_argument_deps={ # dg.AssetKey(["ml_song_summary_info"]), # dg.AssetKey(["dim_clip"]), # dg.AssetKey(["fact_clip_status"]), # }, backfill_policy=dg.BackfillPolicy.multi_run(max_partitions_per_run=7), owners=["team:core-pod"], metadata={ "database": EnvVar("SNOWFLAKE_DB").get_value(), "schema": EnvVar("SNOWFLAKE_SCHEMA").get_value(), "table_name": USER_TOP_SCORE_SONGS_TABLE_NAME, "data_start_date": USER_TOP_SCORE_SONGS_START_DATE.strftime("%Y-%m-%d"), "cluster_by": "[p_date, user_id]", "partition_expr": PartitionExpr.DAILY.value, "transient": True, "sla_minutes": 240, }, ) def user_top_score_songs_snowflake( context: dg.AssetExecutionContext, snowflake: SnowflakeResource, config: UserTopScoreSongsConfig ) -> dg.MaterializeResult: run_id = context.run.run_id logger = dg.get_dagster_logger() jinja_formatter = JinjaSQLFormatter() python_formatter = PythonStringSQLFormatter() # Get partition time window for processing partition_start = context.partition_time_window.start partition_end = context.partition_time_window.end is_multi_partition_range = context.has_partition_key_range # Calculate previous day for deletion/cleanup from datetime import timedelta previous_day = partition_start - timedelta(days=1) fetch_params = { "partition_start_date": partition_start.strftime("%Y-%m-%d"), "partition_end_date": partition_end.strftime("%Y-%m-%d"), "previous_day_date": previous_day.strftime("%Y-%m-%d"), # Add previous day parameter "user_top_score_songs_table_name": config.user_top_score_songs_table_name, } logger.info( f"Processing user_top_score_songs for partition: {partition_start} to {partition_end}" ) logger.info(f"Fetch params: {fetch_params}") with snowflake.get_connection() as conn: cursor = conn.cursor() logger.info(f"Using warehouse {config.warehouse}") warehouse_query = python_formatter.load("src/utils/snowflake/queries/use_warehouse.sql", params={"warehouse": config.warehouse}, logger=logger) cursor.execute(warehouse_query) logger.info( f"Inserting new user top score songs data into {config.user_top_score_songs_table_name}." ) create_query = jinja_formatter.load("src/assets/snowflake/personalization/user_top_score_songs/create.sql", params=fetch_params, logger=logger) cursor.execute(create_query) conn.commit() rows_inserted = cursor.rowcount logger.info(f"Successfully processed partition. Inserted {rows_inserted} rows.") return dg.MaterializeResult( metadata={ "run_id": dg.MetadataValue.text(run_id), "table_name": config.user_top_score_songs_table_name, "partition_time_window_start": dg.MetadataValue.text(partition_start.isoformat()), "partition_time_window_end": dg.MetadataValue.text(partition_end.isoformat()), "dagster/row_count": rows_inserted if not is_multi_partition_range else 0, }, ) @dg.asset( name="user_top_score_songs_ddb", description="Sync user top score songs from Snowflake to DynamoDB via Glue job", group_name=Group.PERSONALIZATION.value, partitions_def=dg.DailyPartitionsDefinition( start_date=USER_TOP_SCORE_SONGS_START_DATE, end_offset=0 ), deps=["user_top_score_songs_snowflake"], backfill_policy=dg.BackfillPolicy.multi_run(max_partitions_per_run=7), owners=["team:core-pod"], metadata={ "sla_minutes": 1440, }, automation_condition=daily_cron_with_eager_historical_backfill_condition ) def user_top_score_songs_ddb(context: dg.AssetExecutionContext) -> dg.MaterializeResult: """ Sync user top score songs from Snowflake to DynamoDB. Steps: 1. Trigger Glue job to read from Snowflake and write to DynamoDB """ run_id = context.run.run_id logger = dg.get_dagster_logger() # Get partition time window for processing partition_start = context.partition_time_window.start partition_end = context.partition_time_window.end logger.info(f"Syncing user_top_score_songs to DynamoDB for partition: {partition_start} to {partition_end}") # Trigger Glue job to sync Snowflake → DynamoDB logger.info("Triggering Glue job to sync Snowflake data to DynamoDB...") try: glue_result = trigger_glue_job( job_name="snowflake_to_ddb_user_top_score_songs_daily", context=context, arguments={ "--partition_date": partition_start.strftime("%Y-%m-%d"), }, poll_interval=60 * 10, # 10 minutes wait_for_completion=True, timeout=3600 * 12, # 12 hours ) logger.info(f"Glue job completed successfully: {glue_result}") glue_job_status = "SUCCESS" glue_job_id = glue_result.get("job_run_id", "N/A") glue_execution_time = glue_result.get("execution_time", 0) except Exception as e: logger.error(f"Glue job failed: {str(e)}") glue_job_status = "FAILED" glue_job_id = "N/A" glue_execution_time = 0 # Re-raise the exception to fail the asset materialization raise e return dg.MaterializeResult( metadata={ "run_id": dg.MetadataValue.text(run_id), "partition_date": dg.MetadataValue.text(partition_start.strftime("%Y-%m-%d")), "glue_job_status": dg.MetadataValue.text(glue_job_status), "glue_job_run_id": dg.MetadataValue.text(glue_job_id), "glue_execution_time_seconds": glue_execution_time, }, )