"""Invite history asset for processing RDS invite and referral 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 INVITE_HISTORY_START_DATE = datetime.strptime('2024-01-01', '%Y-%m-%d') INVITE_HISTORY_TABLE_NAME = "RDS_INVITE_HISTORY" WAREHOUSE = Warehouse.SUNO_PROD_RDS_HOURLY_X_SMALL.value SNOWFLAKE_DB = SnowflakeDB.SUNO_PROD.value SNOWFLAKE_SCHEMA = SnowflakeSchema.PROD.value @dg.asset( name="rds_invite_history", description="Processed invite history table with user invitation and referral data from RDS, triggers Glue export to S3.", group_name=Group.RDS.value, partitions_def=dg.HourlyPartitionsDefinition(start_date=INVITE_HISTORY_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": INVITE_HISTORY_TABLE_NAME, "data_start_date": INVITE_HISTORY_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_invite_history( context: dg.AssetExecutionContext, snowflake: SnowflakeResource ) -> dg.MaterializeResult: """ Materialize the RDS invite history 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, "invite_history_table_name": INVITE_HISTORY_TABLE_NAME, "stage_path": f"@SUNO_DATABASE_EVENTS/bots_invitehistory/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_invitehistory_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": INVITE_HISTORY_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": INVITE_HISTORY_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, } )