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_24H_FAIL_25H, 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 daily_cron_condition warnings.filterwarnings("ignore", category=dg.BetaWarning) ASSET_DIR = Path(__file__).parent AUTH_USER_GROUPS_START_DATE = datetime.strptime('2024-01-01', '%Y-%m-%d') AUTH_USER_GROUPS_TABLE_NAME = "RDS_AUTH_USER_GROUPS" WAREHOUSE = Warehouse.SUNO_PROD_RDS_HOURLY_X_SMALL.value @dg.asset(name="rds_auth_user_groups", description="Auth user groups table", group_name=Group.RDS.value, partitions_def=dg.DailyPartitionsDefinition(start_date=AUTH_USER_GROUPS_START_DATE, end_offset=0), backfill_policy=dg.BackfillPolicy.multi_run(max_partitions_per_run=7), owners=[Team.CORE_POD.value], metadata={"database": SnowflakeDB.SUNO_PROD.value, "schema": SnowflakeSchema.PROD.value, "table_name": AUTH_USER_GROUPS_TABLE_NAME, "partition_expr": PartitionExpr.DAILY.value, "transient": False, "sla_minutes": 60}, automation_condition=daily_cron_condition, freshness_policy=TIME_WINDOW_FRESHNESS_POLICY_WARN_24H_FAIL_25H) def rds_auth_user_groups(context: dg.AssetExecutionContext, snowflake: SnowflakeResource) -> dg.MaterializeResult: logger = dg.get_dagster_logger() 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"), "auth_user_groups_table_name": AUTH_USER_GROUPS_TABLE_NAME, "stage_path": f"@SUNO_DATABASE_EVENTS/auth_usergroups/pdate={partition_start.strftime('%Y-%m-%d')}"} try: glue_result = trigger_glue_job("rds_to_s3_auth_usergroups_daily_all", context, {"--partition_date": partition_start.strftime("%Y-%m-%d")}, 20, True, 1800) glue_job_status, glue_job_id, glue_execution_time = "SUCCESS", glue_result.get("job_run_id", "N/A"), glue_result.get("execution_time", 0) except Exception as e: logger.error(f"Glue job failed: {str(e)}") raise e with snowflake.get_connection() as conn: cursor = conn.cursor() # Step 1: Set warehouse logger.info("Step 1: Use 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) logger.info("Step 2: Merge data into RDS_AUTH_USER_GROUPS") merge_query = load_query(ASSET_DIR / "create.sql", params=fetch_params) log_query(logger, merge_query) cursor.execute(merge_query) # Step 2: Create table logger.info("Step 2: Create Table") create_table_query = load_query(ASSET_DIR / "create.sql", params=fetch_params) log_query(logger, create_table_query) cursor.execute(create_table_query) rows_affected = cursor.rowcount logger.info(f"Successfully created table {AUTH_USER_GROUPS_TABLE_NAME}. Affected {rows_affected} rows.") return dg.MaterializeResult(metadata={"run_id": dg.MetadataValue.text(context.run.run_id), "table_name": AUTH_USER_GROUPS_TABLE_NAME, "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})