"""Genre classification asset for clip embeddings.""" from datetime import datetime import warnings import dagster as dg from src.utils.snowflake.constants import TIME_WINDOW_FRESHNESS_POLICY_WARN_1H_FAIL_2H, Group, Team from src.utils.glue_utils import trigger_glue_job from src.utils.automation_conditions import hourly_cron_with_eager_historical_backfill_condition warnings.filterwarnings("ignore", category=dg.BetaWarning) CLIP_EMBEDDING_START_DATE = datetime.strptime('2024-01-01', '%Y-%m-%d') @dg.asset( name="clip_genre_classification", description="Run genre classification on clip embeddings via Glue job", group_name=Group.RECOMMENDATION.value, partitions_def=dg.HourlyPartitionsDefinition(start_date=CLIP_EMBEDDING_START_DATE, end_offset=0), deps=["rds_clip_embedding"], backfill_policy=dg.BackfillPolicy.multi_run(max_partitions_per_run=24), owners=[Team.CORE_POD.value], metadata={ "sla_minutes": 120, }, automation_condition=hourly_cron_with_eager_historical_backfill_condition, freshness_policy=TIME_WINDOW_FRESHNESS_POLICY_WARN_1H_FAIL_2H, ) def clip_genre_classification(context: dg.AssetExecutionContext) -> dg.MaterializeResult: """ Run genre classification on clip embeddings. This asset depends on rds_clip_embedding and triggers a Glue job to perform genre classification using the embedding data. Steps: 1. Trigger Glue job for recommendation genre classification """ 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"Running genre classification for partition: {partition_start} to {partition_end}") # Trigger Glue job for genre classification logger.info("Triggering Glue job for genre classification...") try: glue_result = trigger_glue_job( job_name="recommendation_genre_classification", context=context, arguments={ "--partition_date": partition_start.strftime("%Y-%m-%d"), "--partition_hour": str(partition_start.hour), }, 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_date": dg.MetadataValue.text(partition_start.strftime("%Y-%m-%d")), "partition_hour": partition_start.hour, "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, }, )