"""Asset checks for ddb_item_info to validate data quality.""" import dagster as dg from dagster_snowflake import SnowflakeResource from src.assets.snowflake.raw_table.dynamodb.item_info.assets import DDB_ITEM_INFO_TABLE_NAME from src.utils.snowflake.partition_utils import parse_hourly_partition_key from src.utils.snowflake.constants import Warehouse # Constants for validation MIN_ROW_COUNT = 0 # Minimum row count per partition (0 allows empty partitions) MAX_DECREASE_RATE = 0.70 # Maximum allowed decrease rate (70%) @dg.asset_check(asset="ddb_item_info", blocking=False) def check_minimum_row_count(context: dg.AssetCheckExecutionContext, snowflake: SnowflakeResource) -> dg.AssetCheckResult: """Ensure minimum row count threshold is met.""" partition_key = context.op_execution_context.partition_key p_date, p_hour = parse_hourly_partition_key(partition_key) with snowflake.get_connection() as conn: cursor = conn.cursor() # Use the appropriate warehouse cursor.execute(f"USE WAREHOUSE {Warehouse.DYNAMODB_EVENTS_SMALL.value}") query = f""" SELECT COUNT(*) as row_count FROM {DDB_ITEM_INFO_TABLE_NAME} WHERE p_date = '{p_date}' and p_hour = {p_hour} """ result = cursor.execute(query).fetchone() row_count = result[0] if result else 0 passed = row_count >= MIN_ROW_COUNT return dg.AssetCheckResult( passed=passed, metadata={ "partition_key": partition_key, "p_date": p_date, "p_hour": p_hour, "row_count": row_count, "min_threshold": MIN_ROW_COUNT, }, description=f"Row count: {row_count:,} (threshold: {MIN_ROW_COUNT:,})" ) @dg.asset_check(asset="ddb_item_info", blocking=False) def check_id_type_uniqueness(context: dg.AssetCheckExecutionContext, snowflake: SnowflakeResource) -> dg.AssetCheckResult: """Verify that (id, type) combination is unique within each partition.""" partition_key = context.op_execution_context.partition_key p_date, p_hour = parse_hourly_partition_key(partition_key) with snowflake.get_connection() as conn: cursor = conn.cursor() # Use the appropriate warehouse cursor.execute(f"USE WAREHOUSE {Warehouse.DYNAMODB_EVENTS_SMALL.value}") # Check for duplicate (id, type) combinations query = f""" SELECT COUNT(*) as total_rows, COUNT(DISTINCT id, type) as unique_combinations, COUNT(*) - COUNT(DISTINCT id, type) as duplicate_count FROM {DDB_ITEM_INFO_TABLE_NAME} WHERE p_date = '{p_date}' and p_hour = {p_hour} """ result = cursor.execute(query).fetchone() if result: total_rows = result[0] unique_combinations = result[1] duplicate_count = result[2] else: total_rows = 0 unique_combinations = 0 duplicate_count = 0 passed = duplicate_count == 0 return dg.AssetCheckResult( passed=passed, metadata={ "partition_key": partition_key, "p_date": p_date, "p_hour": p_hour, "total_rows": total_rows, "unique_combinations": unique_combinations, "duplicate_count": duplicate_count, }, description=f"Found {duplicate_count} duplicate (id, type) combinations out of {total_rows} total rows" ) @dg.asset_check(asset="ddb_item_info", blocking=False) def check_day_over_day_volume(context: dg.AssetCheckExecutionContext, snowflake: SnowflakeResource) -> dg.AssetCheckResult: """Check that today's data volume is growing or decreasing by less than 20%.""" partition_key = context.op_execution_context.partition_key p_date, p_hour = parse_hourly_partition_key(partition_key) with snowflake.get_connection() as conn: cursor = conn.cursor() # Use the appropriate warehouse cursor.execute(f"USE WAREHOUSE {Warehouse.DYNAMODB_EVENTS_SMALL.value}") # Get both today's and yesterday's count in one query using conditional aggregation query = f""" SELECT SUM(CASE WHEN p_date = '{p_date}' THEN 1 ELSE 0 END) as today_count, SUM(CASE WHEN p_date = DATEADD(day, -1, '{p_date}'::date) THEN 1 ELSE 0 END) as yesterday_count FROM {DDB_ITEM_INFO_TABLE_NAME} WHERE p_date IN ('{p_date}', DATEADD(day, -1, '{p_date}'::date)) AND p_hour = {p_hour} """ result = cursor.execute(query).fetchone() today_count = result[0] if result and result[0] else 0 yesterday_count = result[1] if result and result[1] else 0 # Check volume change if yesterday_count == 0: # If no data yesterday, skip this check return dg.AssetCheckResult( passed=True, description=f"No data for yesterday's partition ({p_date} {p_hour}:00), skipping volume comparison.", metadata={ "partition_key": partition_key, "today_count": today_count, "yesterday_count": yesterday_count, } ) # Calculate change change = today_count - yesterday_count change_percentage = (change / yesterday_count) * 100 if yesterday_count > 0 else 0 # Pass if: # 1. Today > yesterday (growth is always good) # 2. Today <= yesterday BUT decrease is less than 20% if today_count > yesterday_count: passed = True reason = "growing" else: decrease_rate = abs(change) / yesterday_count if decrease_rate < MAX_DECREASE_RATE: passed = True reason = f"decrease within acceptable range ({decrease_rate*100:.2f}% < {MAX_DECREASE_RATE*100:.0f}%)" else: passed = False reason = f"decrease exceeds threshold ({decrease_rate*100:.2f}% >= {MAX_DECREASE_RATE*100:.0f}%)" return dg.AssetCheckResult( passed=passed, description=f"Today: {today_count}, Yesterday: {yesterday_count}, Change: {change:+d} ({change_percentage:+.2f}%) - {reason}", metadata={ "partition_key": partition_key, "today_count": today_count, "yesterday_count": yesterday_count, "change": change, "change_percentage": f"{change_percentage:+.2f}%", "status": reason, } ) # Export all checks as a list asset_checks = [ check_minimum_row_count, check_id_type_uniqueness, check_day_over_day_volume, ] __all__ = ["asset_checks"]