""" Glue Job Orchestration for Self Listening Clips """ from datetime import datetime import warnings import dagster as dg from src.utils.snowflake.constants import Group, Team from src.utils.glue_utils import trigger_glue_job from src.utils.automation_conditions import daily_cron_with_eager_historical_backfill_condition warnings.filterwarnings("ignore", category=dg.BetaWarning) # use same start date as fact_play SELF_LISTENING_CLIPS_PARTITION_START_DATE = datetime.strptime('2024-06-01', '%Y-%m-%d') @dg.asset( name="self_listening_clips", description="Write top 25 self listening clips from Snowflake to Redis via Glue job", group_name=Group.LISTENING.value, partitions_def=dg.TimeWindowPartitionsDefinition( start=SELF_LISTENING_CLIPS_PARTITION_START_DATE, cron_schedule="0 0 * * *", # Create partitions every day at 00:00 fmt="%Y-%m-%d-%H:%M", end_offset=0, ), deps=[ dg.AssetDep("fact_play", partition_mapping=dg.TimeWindowPartitionMapping()), dg.AssetDep("clip", partition_mapping=dg.TimeWindowPartitionMapping()), ], backfill_policy=dg.BackfillPolicy.multi_run(max_partitions_per_run=7), owners=[Team.CORE_POD.value], metadata={ "sla_minutes": 240, }, automation_condition=daily_cron_with_eager_historical_backfill_condition, ) def self_listening_clips(context: dg.AssetExecutionContext) -> dg.MaterializeResult: """ Write top 25 self listening clips from Snowflake to Redis. This asset depends on fact_play and clip, and triggers a Glue job to query self-listening data from Snowflake and write it to Redis. Steps: 1. Trigger Glue job to read from Snowflake and write to Redis """ 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 self listening clips to Redis for partition: {partition_start} to {partition_end}") # Trigger Glue job to sync Snowflake → Redis logger.info("Triggering Glue job to sync self listening clips from Snowflake to Redis...") try: glue_result = trigger_glue_job( job_name="snowflake_to_redis_self_listening_clips", context=context, arguments=None, # Glue job queries last 120 days without partition args poll_interval=60, # 1 minute wait_for_completion=True, timeout=3600 * 2, # 2 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_time_window_start": dg.MetadataValue.text(partition_start.isoformat()), "partition_time_window_end": dg.MetadataValue.text(partition_end.isoformat()), "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, }, )