"""Auth user asset for processing RDS auth user 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 hourly_cron_with_eager_historical_backfill_condition warnings.filterwarnings("ignore", category=dg.BetaWarning) # Get directory of this file for relative SQL file loading ASSET_DIR = Path(__file__).parent AUTH_USER_START_DATE = datetime.strptime('2024-01-01', '%Y-%m-%d') AUTH_USER_TABLE_NAME = "RDS_AUTH_USER" DELETED_AUTH_USER_TABLE_NAME = "DELETED_RDS_AUTH_USER" WAREHOUSE = Warehouse.SUNO_PROD_RDS_HOURLY_X_SMALL.value SNOWFLAKE_DB = SnowflakeDB.SUNO_PROD.value SNOWFLAKE_SCHEMA = SnowflakeSchema.PROD.value @dg.asset( name="rds_auth_user", description="Auth user table (hourly partitioned) containing authentication user data from RDS", group_name=Group.RDS.value, partitions_def=dg.HourlyPartitionsDefinition(start_date=AUTH_USER_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": AUTH_USER_TABLE_NAME, "cluster_by": "[p_date, p_hour]", "partition_expr": PartitionExpr.HOURLY.value, "transient": True, "sla_minutes": 60, }, automation_condition=hourly_cron_with_eager_historical_backfill_condition, freshness_policy=TIME_WINDOW_FRESHNESS_POLICY_WARN_1H_FAIL_2H, ) def rds_auth_user( context: dg.AssetExecutionContext, snowflake: SnowflakeResource ) -> dg.MaterializeResult: """ Process RDS auth user data for a specific hourly partition. This asset: 1. Triggers AWS 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 information snowflake: Snowflake resource for database operations Returns: MaterializeResult with metadata about the operation """ run_id = context.run.run_id logger = dg.get_dagster_logger() 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, "auth_user_table_name": AUTH_USER_TABLE_NAME, "stage_path": f"@SUNO_DATABASE_EVENTS/auth_user/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_auth_user_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 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) 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 with snowflake.get_connection() as conn: cursor = conn.cursor() # Step 2: Set warehouse logger.info("Step 2: Use Warehouse") cursor.execute( load_query( "src/utils/snowflake/queries/use_warehouse.sql", params={"warehouse": WAREHOUSE} ) ) # Step 4: Delete existing partition data logger.info("Step 4: Delete Hourly Partition") cursor.execute( load_query( "src/utils/snowflake/queries/delete_hourly_partitions.sql", params={**fetch_params, "delete_partition_table_name": AUTH_USER_TABLE_NAME} ) ) # Step 5: Insert data from S3 logger.info("Step 5: Insert Data") 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 inserted data into {AUTH_USER_TABLE_NAME}. Affected {rows_affected} rows.") return dg.MaterializeResult( metadata={ "run_id": dg.MetadataValue.text(run_id), "table_name": AUTH_USER_TABLE_NAME, "partition_time_window_start": dg.MetadataValue.text(partition_start.isoformat()), "partition_time_window_end": dg.MetadataValue.text(partition_end.isoformat()), "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, } ) @dg.asset( name="deleted_rds_auth_user", description="Deleted auth user table containing auth users that don't exist in discord_info (daily partitioned)", group_name=Group.RDS.value, partitions_def=dg.DailyPartitionsDefinition(start_date=AUTH_USER_START_DATE, end_offset=0), backfill_policy=dg.BackfillPolicy.multi_run(max_partitions_per_run=7), deps=[ dg.AssetDep("rds_discord_info", partition_mapping=dg.TimeWindowPartitionMapping(start_offset=-1, end_offset=0)), ], owners=[Team.CORE_POD.value], metadata={ "database": SNOWFLAKE_DB, "schema": SNOWFLAKE_SCHEMA, "table_name": DELETED_AUTH_USER_TABLE_NAME, "cluster_by": "[p_date, p_hour]", "partition_expr": PartitionExpr.DAILY.value, "transient": True, "sla_minutes": 60, }, automation_condition=dg.AutomationCondition.on_cron(cron_schedule="0 1 * * *", cron_timezone="UTC"), ) def deleted_rds_auth_user( context: dg.AssetExecutionContext, snowflake: SnowflakeResource ) -> dg.MaterializeResult: """ Move auth users that don't exist in discord_info to deleted table. This asset runs daily and: 1. Creates the deleted table if it doesn't exist 2. Finds auth_user records for the entire day where ID is NOT IN discord_info 3. Inserts those records into the deleted table 4. Deletes those records from the main auth_user table Args: context: Dagster execution context with partition information snowflake: Snowflake resource for database operations Returns: MaterializeResult with metadata about the operation """ run_id = context.run.run_id logger = dg.get_dagster_logger() partition_start = context.partition_time_window.start partition_end = context.partition_time_window.end partition_date = partition_start.strftime("%Y-%m-%d") logger.info(f"Processing daily partition: {partition_date}") # Prepare SQL parameters fetch_params = { "partition_start_date": partition_date, "partition_end_date": partition_date, "auth_user_table_name": AUTH_USER_TABLE_NAME, "deleted_auth_user_table_name": DELETED_AUTH_USER_TABLE_NAME, } rows_moved = 0 rows_deleted = 0 with snowflake.get_connection() as conn: cursor = conn.cursor() # Step 1: Set warehouse logger.info("Step 1: Use Warehouse") cursor.execute( load_query( "src/utils/snowflake/queries/use_warehouse.sql", params={"warehouse": WAREHOUSE} ) ) # Step 2: Create deleted table if not exists logger.info("Step 2: Create Deleted Table") cursor.execute( load_query( ASSET_DIR / "deleted_table.sql", params=fetch_params ) ) # Step 3: Delete existing partition data from deleted table (for the entire day) logger.info("Step 3: Delete Existing Partition from Deleted Table") delete_partition_query = f""" DELETE FROM {SNOWFLAKE_DB}.{SNOWFLAKE_SCHEMA}.{DELETED_AUTH_USER_TABLE_NAME} WHERE P_DATE = '{partition_date}' """ log_query(logger, delete_partition_query) cursor.execute(delete_partition_query) # Step 4: Find and insert auth users not in discord_info (for entire day) logger.info("Step 4: Insert auth users not in discord_info into deleted table") insert_query = f""" INSERT INTO {SNOWFLAKE_DB}.{SNOWFLAKE_SCHEMA}.{DELETED_AUTH_USER_TABLE_NAME} (ID, USERNAME, DATE_JOINED, P_DATE, P_HOUR, EMAIL, IS_STAFF) SELECT au.ID, au.USERNAME, au.DATE_JOINED, au.P_DATE, au.P_HOUR, au.EMAIL, au.IS_STAFF FROM {SNOWFLAKE_DB}.{SNOWFLAKE_SCHEMA}.{AUTH_USER_TABLE_NAME} au WHERE au.P_DATE = '{partition_date}' AND au.ID NOT IN ( SELECT DISTINCT user_id FROM {SNOWFLAKE_DB}.{SNOWFLAKE_SCHEMA}.RDS_DISCORD_INFO WHERE user_id IS NOT NULL ) """ log_query(logger, insert_query) cursor.execute(insert_query) rows_moved = cursor.rowcount logger.info(f"Moved {rows_moved} records to deleted table") # Step 5: Delete those records from main auth_user table (for entire day) logger.info("Step 5: Delete moved records from main auth_user table") delete_query = f""" DELETE FROM {SNOWFLAKE_DB}.{SNOWFLAKE_SCHEMA}.{AUTH_USER_TABLE_NAME} WHERE P_DATE = '{partition_date}' AND ID NOT IN ( SELECT DISTINCT user_id FROM {SNOWFLAKE_DB}.{SNOWFLAKE_SCHEMA}.RDS_DISCORD_INFO WHERE user_id IS NOT NULL ) """ log_query(logger, delete_query) cursor.execute(delete_query) rows_deleted = cursor.rowcount logger.info(f"Deleted {rows_deleted} records from main table") logger.info(f"Successfully processed partition. Moved {rows_moved} rows to deleted table.") return dg.MaterializeResult( metadata={ "run_id": dg.MetadataValue.text(run_id), "table_name": DELETED_AUTH_USER_TABLE_NAME, "partition_time_window_start": dg.MetadataValue.text(partition_start.isoformat()), "partition_time_window_end": dg.MetadataValue.text(partition_end.isoformat()), "rows_moved_to_deleted": rows_moved, "rows_deleted_from_main": rows_deleted, } )