"""Orpheus asset for processing DynamoDB orpheus data and triggering Glue export.""" from datetime import datetime import warnings import dagster as dg from dagster_snowflake import SnowflakeResource from src.utils.automation_conditions import hourly_cron_condition from src.utils.snowflake.constants import TIME_WINDOW_FRESHNESS_POLICY_WARN_1H_FAIL_2H, PartitionExpr, Team, Group from src.utils.glue_utils import trigger_glue_job warnings.filterwarnings("ignore", category=dg.BetaWarning) ORPHEUS_PROJECT_START_DATE = datetime.strptime('2025-01-01:00', '%Y-%m-%d:%H') ORPHEUS_PROJECT_TABLE_NAME = "RAW_S3_DDB_ORPHEUS" @dg.asset( name="ddb_orpheus", description="DynamoDB S3 dump of orpheus data with JSON metadata", group_name=Group.DYNAMODB.value, partitions_def=dg.HourlyPartitionsDefinition(start_date=ORPHEUS_PROJECT_START_DATE, end_offset=0), backfill_policy=dg.BackfillPolicy.multi_run(max_partitions_per_run=24), automation_condition=hourly_cron_condition, owners=[Team.CORE_POD.value], metadata={ "table_name": ORPHEUS_PROJECT_TABLE_NAME, "data_start_date": ORPHEUS_PROJECT_START_DATE.strftime("%Y-%m-%d"), "cluster_by": "[p_date, p_hour]", "partition_expr": PartitionExpr.HOURLY.value, "sla_minutes": 60, }, freshness_policy=TIME_WINDOW_FRESHNESS_POLICY_WARN_1H_FAIL_2H, ) def ddb_orpheus( context: dg.AssetExecutionContext, ) -> dg.MaterializeResult: """ Process orpheus data from DDB and trigger Glue export. Steps: 1. Trigger Glue job to export data to S3 """ 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 fetch_params = { "partition_start_date": partition_start.strftime("%Y-%m-%d"), "partition_end_date": partition_end.strftime("%Y-%m-%d"), "partition_start_hour": partition_start.hour, "partition_end_hour": partition_end.hour, "orpheus_table_name": ORPHEUS_PROJECT_TABLE_NAME } logger.info(f"Processing orpheus for partition: {partition_start} to {partition_end}") logger.info(f"Fetch params: {fetch_params}") try: glue_result = trigger_glue_job( job_name="ddb_to_s3_orpheus_hourly_insert", context=context, arguments={ "--partition_date": partition_start.strftime("%Y-%m-%d"), "--partition_hour": str(partition_start.hour), }, poll_interval=20, wait_for_completion=True, timeout=1800, # 30 minutes ) 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), "table_name": ORPHEUS_PROJECT_TABLE_NAME, "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, }, )