import dagster as dg from dagster_snowflake import SnowflakeResource from src.utils.snowflake.constants import Warehouse from src.utils.snowflake.query import load_query from src.utils.snowflake.partition_utils import parse_hourly_partition_key PROJECT_CLIP_TABLE_NAME = "RDS_PROJECT_CLIP" WAREHOUSE = Warehouse.SUNO_PROD_RDS_HOURLY_X_SMALL.value @dg.asset_check(asset="rds_project_clip", description="Verify ID uniqueness within partition", blocking=False) def check_project_clip_id_uniqueness(context: dg.AssetCheckExecutionContext, snowflake: SnowflakeResource) -> dg.AssetCheckResult: """ Checks for ID uniqueness within the current hourly partition of the RDS_PROJECT_CLIP table. """ partition_key = context.op_execution_context.partition_key # Parse partition key (format: YYYY-MM-DD-HH:00) # Example: "2025-10-03-14:00" partition_parts = partition_key.split("-") p_date = f"{partition_parts[0]}-{partition_parts[1]}-{partition_parts[2]}" p_hour = int(partition_parts[3].split(":")[0]) with snowflake.get_connection() as conn: cursor = conn.cursor() cursor.execute(load_query("src/utils/snowflake/queries/use_warehouse.sql", params={"warehouse": WAREHOUSE})) query = f""" SELECT COUNT(*), COUNT(DISTINCT ID), COUNT(*) - COUNT(DISTINCT ID) FROM {PROJECT_CLIP_TABLE_NAME} WHERE P_DATE = '{p_date}' AND P_HOUR = {p_hour} """ cursor.execute(query) 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, "partition_key": partition_key, } ) @dg.asset_check(asset="rds_project_clip", description="Verify row count > 0 for the partition", blocking=False) def check_project_clip_row_count(context: dg.AssetCheckExecutionContext, snowflake: SnowflakeResource) -> dg.AssetCheckResult: """ Checks that the row count for the current hourly partition of the RDS_PROJECT_CLIP table is greater than 0. """ partition_key = context.op_execution_context.partition_key # Parse partition key (format: YYYY-MM-DD-HH:00) # Example: "2025-10-03-14:00" partition_parts = partition_key.split("-") p_date = f"{partition_parts[0]}-{partition_parts[1]}-{partition_parts[2]}" p_hour = int(partition_parts[3].split(":")[0]) with snowflake.get_connection() as conn: cursor = conn.cursor() cursor.execute(load_query("src/utils/snowflake/queries/use_warehouse.sql", params={"warehouse": WAREHOUSE})) query = f""" SELECT COUNT(*) FROM {PROJECT_CLIP_TABLE_NAME} WHERE P_DATE = '{p_date}' AND P_HOUR = {p_hour} """ cursor.execute(query) row_count = cursor.fetchone()[0] return dg.AssetCheckResult( passed=row_count > 0, description=f"Table has {row_count} rows for partition {partition_key}", metadata={ "row_count": row_count, "partition_key": partition_key, } ) asset_checks = [ check_project_clip_id_uniqueness, check_project_clip_row_count, ] __all__ = ["asset_checks"]