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_contest", description="Verify ID uniqueness", blocking=False) def check_contest_id_uniqueness(context: dg.AssetCheckExecutionContext, snowflake: SnowflakeResource) -> dg.AssetCheckResult: p_date = context.op_execution_context.partition_key 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(f"SELECT COUNT(*), COUNT(DISTINCT ID), COUNT(*) - COUNT(DISTINCT ID) FROM SUNO_PROD.PROD.RDS_CONTEST WHERE p_date = '{p_date}'") 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, "duplicates": dupes}) asset_checks = [check_contest_id_uniqueness] __all__ = ["asset_checks"]