"""Usage plan discount offer asset for processing RDS usage plan discount offer 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) ASSET_DIR = Path(__file__).parent USAGE_PLAN_DISCOUNT_OFFER_START_DATE = datetime.strptime('2024-01-01', '%Y-%m-%d') USAGE_PLAN_DISCOUNT_OFFER_TABLE_NAME = "RDS_USAGE_PLAN_DISCOUNT_OFFER" WAREHOUSE = Warehouse.SUNO_PROD_RDS_HOURLY_X_SMALL.value SNOWFLAKE_DB = SnowflakeDB.SUNO_PROD.value SNOWFLAKE_SCHEMA = SnowflakeSchema.PROD.value @dg.asset( name="rds_usage_plan_discount_offer", description="Processed usage plan discount offer table with data from RDS, triggers Glue export to S3.", group_name=Group.RDS.value, partitions_def=dg.HourlyPartitionsDefinition(start_date=USAGE_PLAN_DISCOUNT_OFFER_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": USAGE_PLAN_DISCOUNT_OFFER_TABLE_NAME, "data_start_date": USAGE_PLAN_DISCOUNT_OFFER_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_usage_plan_discount_offer(context: dg.AssetExecutionContext, snowflake: SnowflakeResource) -> dg.MaterializeResult: run_id = context.run.run_id logger = dg.get_dagster_logger() partition_start = context.partition_time_window.start fetch_params = { "partition_start_date": partition_start.strftime("%Y-%m-%d"), "partition_end_date": context.partition_time_window.end.strftime("%Y-%m-%d"), "partition_start_hour": partition_start.hour, "partition_end_hour": context.partition_time_window.end.hour, "usage_plan_discount_offer_table_name": USAGE_PLAN_DISCOUNT_OFFER_TABLE_NAME, "stage_path": f"@SUNO_DATABASE_EVENTS/billing_usageplandiscountoffer/pdate={partition_start.strftime('%Y-%m-%d')}/phour={partition_start.strftime('%H')}", } rows_affected = 0 try: glue_result = trigger_glue_job("rds_to_s3_billing_usageplandiscountoffer_hourly_upsert", context, {"--partition_date": partition_start.strftime("%Y-%m-%d"), "--partition_hour": str(partition_start.hour)}, 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)}") glue_job_status, glue_job_id, glue_execution_time = "FAILED", "N/A", 0 raise e with snowflake.get_connection() as conn: cursor = conn.cursor() cursor.execute(load_query("src/utils/snowflake/queries/use_warehouse.sql", params={"warehouse": WAREHOUSE})) cursor.execute(load_query("src/utils/snowflake/queries/delete_hourly_partitions.sql", params={**fetch_params, "delete_partition_table_name": USAGE_PLAN_DISCOUNT_OFFER_TABLE_NAME})) cursor.execute(load_query(ASSET_DIR / "upsert.sql", params=fetch_params)) rows_affected = cursor.rowcount return dg.MaterializeResult(metadata={"run_id": dg.MetadataValue.text(run_id), "table_name": USAGE_PLAN_DISCOUNT_OFFER_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})