"""Asset checks for rds_usage_plan data quality.""" import dagster as dg from dagster_snowflake import SnowflakeResource from src.utils.snowflake.constants import Warehouse from src.utils.snowflake.query import load_query @dg.asset_check( asset="rds_usage_plan", description="Verify ID uniqueness", blocking=False ) def check_usage_plan_id_uniqueness( context: dg.AssetCheckExecutionContext, snowflake: SnowflakeResource ) -> dg.AssetCheckResult: """Check that all IDs in the usage plan table are unique.""" with snowflake.get_connection() as conn: cursor = conn.cursor() cursor.execute(load_query( "src/utils/snowflake/queries/use_warehouse.sql", params={"warehouse": Warehouse.SUNO_PROD_RDS_HOURLY_X_SMALL.value} )) cursor.execute( "SELECT COUNT(*), COUNT(DISTINCT ID), COUNT(*) - COUNT(DISTINCT ID) " "FROM SUNO_PROD.PROD.RDS_USAGE_PLAN" ) total, distinct, dupes = cursor.fetchone() or (0, 0, 0) return dg.AssetCheckResult( passed=dupes == 0, description=f"{'Unique' if dupes == 0 else f'{dupes} duplicates'}", metadata={ "total_rows": total, "distinct_ids": distinct, "duplicates": dupes } ) @dg.asset_check( asset="rds_usage_plan", description="Verify row count > 0", blocking=False ) def check_usage_plan_row_count( context: dg.AssetCheckExecutionContext, snowflake: SnowflakeResource ) -> dg.AssetCheckResult: """Check that the usage plan table has data.""" with snowflake.get_connection() as conn: cursor = conn.cursor() cursor.execute(load_query( "src/utils/snowflake/queries/use_warehouse.sql", params={"warehouse": Warehouse.SUNO_PROD_RDS_HOURLY_X_SMALL.value} )) cursor.execute("SELECT COUNT(*) FROM SUNO_PROD.PROD.RDS_USAGE_PLAN") row_count = cursor.fetchone()[0] return dg.AssetCheckResult( passed=row_count > 0, description=f"Table has {row_count} rows", metadata={"row_count": row_count} ) # Export all asset checks asset_checks = [ check_usage_plan_id_uniqueness, check_usage_plan_row_count, ] __all__ = ["asset_checks"]