"""Asset checks for RDS profile follow 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 def parse_daily_partition_key(partition_key: str) -> str: """Parse daily partition key. Args: partition_key: Partition key in format 'YYYY-MM-DD' Returns: p_date string """ return partition_key @dg.asset_check( asset="rds_profile_follow", description="Verify ID uniqueness within each partition", blocking=False ) def check_profile_follow_id_uniqueness( context: dg.AssetCheckExecutionContext, snowflake: SnowflakeResource ) -> dg.AssetCheckResult: """Check that IDs are unique within the partition.""" partition_key = context.op_execution_context.partition_key p_date = parse_daily_partition_key(partition_key) with snowflake.get_connection() as conn: cursor = conn.cursor() # Set warehouse cursor.execute( load_query( "src/utils/snowflake/queries/use_warehouse.sql", params={"warehouse": Warehouse.SUNO_PROD_RDS_HOURLY_X_SMALL.value} ) ) # Check for duplicate IDs in this partition cursor.execute( f""" SELECT COUNT(*), COUNT(DISTINCT ID), COUNT(*) - COUNT(DISTINCT ID) FROM SUNO_PROD.PROD.RDS_PROFILE_FOLLOW 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'} in partition {partition_key}", metadata={ "partition_key": partition_key, "p_date": p_date, "total_rows": total, "duplicates": dupes, } ) @dg.asset_check( asset="rds_profile_follow", description="Verify row count is greater than 0", blocking=False ) def check_profile_follow_row_count( context: dg.AssetCheckExecutionContext, snowflake: SnowflakeResource ) -> dg.AssetCheckResult: """Check that the partition has at least one row.""" partition_key = context.op_execution_context.partition_key p_date = parse_daily_partition_key(partition_key) with snowflake.get_connection() as conn: cursor = conn.cursor() # Set warehouse cursor.execute( load_query( "src/utils/snowflake/queries/use_warehouse.sql", params={"warehouse": Warehouse.SUNO_PROD_RDS_HOURLY_X_SMALL.value} ) ) # Get row count for this partition cursor.execute( f""" SELECT COUNT(*) FROM SUNO_PROD.PROD.RDS_PROFILE_FOLLOW WHERE p_date = '{p_date}' """ ) row_count = cursor.fetchone()[0] return dg.AssetCheckResult( passed=row_count > 0, description=f"Partition has {row_count} rows", metadata={ "partition_key": partition_key, "p_date": p_date, "row_count": row_count, } ) # Export all checks asset_checks = [ check_profile_follow_id_uniqueness, check_profile_follow_row_count, ] __all__ = ["asset_checks"]