"""Grouped user notification asset for processing RDS grouped notification data and triggering Glue export.""" from datetime import datetime import warnings from pathlib import Path import dagster as dg from dagster_snowflake import SnowflakeResource from src.utils.snowflake.constants import TIME_WINDOW_FRESHNESS_POLICY_WARN_1H_FAIL_2H, Group, PartitionExpr, Warehouse, SnowflakeDB, Team, SnowflakeSchema from src.utils.snowflake.logger import log_query from src.utils.snowflake.query import load_query from src.utils.glue_utils import trigger_glue_job from src.utils.automation_conditions import rds_hourly_cron_condition warnings.filterwarnings("ignore", category=dg.BetaWarning) # Get directory of this file for relative SQL file loading ASSET_DIR = Path(__file__).parent GROUPED_USER_NOTIFICATION_START_DATE = datetime.strptime('2024-01-01', '%Y-%m-%d') GROUPED_USER_NOTIFICATION_TABLE_NAME = "RDS_GROUPED_USER_NOTIFICATION" WAREHOUSE = Warehouse.SUNO_PROD_RDS_HOURLY_X_SMALL.value SNOWFLAKE_DB = SnowflakeDB.SUNO_PROD.value SNOWFLAKE_SCHEMA = SnowflakeSchema.PROD.value @dg.asset( name="rds_grouped_user_notification", description="Processed grouped user notification table with aggregated notification data from RDS, triggers Glue export to S3.", group_name=Group.RDS.value, partitions_def=dg.HourlyPartitionsDefinition(start_date=GROUPED_USER_NOTIFICATION_START_DATE, end_offset=0), backfill_policy=dg.BackfillPolicy.multi_run(max_partitions_per_run=24), owners=[Team.CORE_POD.value], metadata={ "database": SNOWFLAKE_DB, "schema": SNOWFLAKE_SCHEMA, "table_name": GROUPED_USER_NOTIFICATION_TABLE_NAME, "data_start_date": GROUPED_USER_NOTIFICATION_START_DATE.strftime("%Y-%m-%d"), "cluster_by": "[p_date, p_hour]", "partition_expr": PartitionExpr.HOURLY.value, "transient": True, "sla_minutes": 60, }, automation_condition=rds_hourly_cron_condition, freshness_policy=TIME_WINDOW_FRESHNESS_POLICY_WARN_1H_FAIL_2H, ) def rds_grouped_user_notification( context: dg.AssetExecutionContext, snowflake: SnowflakeResource ) -> dg.MaterializeResult: """ Materialize the RDS grouped user notification table. This asset: 1. Triggers a Glue job to export data from RDS to S3 2. Creates the Snowflake table if it doesn't exist 3. Deletes existing data for the partition 4. Loads new data from S3 into Snowflake Args: context: Dagster execution context with partition info snowflake: Snowflake resource for database connections Returns: MaterializeResult with metadata about the operation """ run_id = context.run.run_id logger = dg.get_dagster_logger() # Get partition time window partition_start = context.partition_time_window.start partition_end = context.partition_time_window.end logger.info(f"Processing partition: {partition_start} to {partition_end}") # Prepare SQL parameters 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, "grouped_user_notification_table_name": GROUPED_USER_NOTIFICATION_TABLE_NAME, "stage_path": f"@SUNO_DATABASE_EVENTS/bots_groupedusernotification/pdate={partition_start.strftime('%Y-%m-%d')}/phour={partition_start.strftime('%H')}", } rows_affected = 0 try: # Step 1: Trigger Glue job to export from RDS to S3 logger.info("Step 1: Triggering Glue job to export data from RDS to S3") glue_result = trigger_glue_job( job_name="rds_to_s3_bots_groupedusernotification_hourly_upsert", context=context, arguments={ "--partition_date": partition_start.strftime("%Y-%m-%d"), "--partition_hour": str(partition_start.hour) }, wait_for_completion=True, poll_interval=20, timeout=1800 # 30 minutes timeout ) glue_job_status = "SUCCESS" glue_job_id = glue_result.get("job_run_id", "N/A") glue_execution_time = glue_result.get("execution_time", 0) logger.info(f"Glue job completed successfully. Job Run ID: {glue_job_id}, Execution Time: {glue_execution_time}s") 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 raise e # Step 2: Process in Snowflake with snowflake.get_connection() as conn: cursor = conn.cursor() # Step 2a: Use warehouse logger.info("Step 2a: Setting warehouse") use_warehouse_query = load_query( "src/utils/snowflake/queries/use_warehouse.sql", params={"warehouse": WAREHOUSE} ) log_query(logger, use_warehouse_query) cursor.execute(use_warehouse_query) # Step 2c: Delete existing hourly partition data logger.info("Step 2c: Deleting existing partition data") delete_partition_query = load_query( "src/utils/snowflake/queries/delete_hourly_partitions.sql", params={ **fetch_params, "delete_partition_table_name": GROUPED_USER_NOTIFICATION_TABLE_NAME } ) log_query(logger, delete_partition_query) cursor.execute(delete_partition_query) # Step 2d: Upsert data from S3 logger.info("Step 2d: Upserting data from S3") upsert_query = load_query(ASSET_DIR / "upsert.sql", params=fetch_params) log_query(logger, upsert_query) cursor.execute(upsert_query) rows_affected = cursor.rowcount logger.info(f"Successfully processed {rows_affected} rows for partition {partition_start}") return dg.MaterializeResult( metadata={ "run_id": dg.MetadataValue.text(run_id), "table_name": GROUPED_USER_NOTIFICATION_TABLE_NAME, "dagster/row_count": rows_affected, "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, } )