from snowflake.snowpark.session import Session from datetime import datetime, timedelta def validate(session: Session, CUR_PDATE: str, CUR_PHOUR: int, PRE_PDATE: str, PRE_PHOUR: int): query_a = """ SELECT COUNT(*) AS num FROM ML_SONG_SUMMARY_INFO WHERE p_date = '%s' AND p_hour = %s """ query_b = """ SELECT COUNT(*) AS num FROM ML_SONG_SUMMARY_INFO WHERE p_date = '%s' AND p_hour = %s """ result_a = session.sql(query_a % (CUR_PDATE, CUR_PHOUR)).collect()[0][0] result_b = session.sql(query_b % (PRE_PDATE, PRE_PHOUR)).collect()[0][0] # Combine pre_date and pre_hour into a datetime object pre_datetime = datetime.strptime(PRE_PDATE, "%Y-%m-%d") + timedelta(hours=PRE_PHOUR) # Subtract one hour, temporary to keep one more hour data one_hour_earlier = pre_datetime - timedelta(hours=2) # Extract the date and hour earlier_date = one_hour_earlier.strftime("%Y-%m-%d") earlier_hour = one_hour_earlier.hour # if the current data is greater than the previous data, delete the previous data # otherwise, send a slack alert if result_a > result_b and result_a > 0: delete_query = """ DELETE FROM ML_SONG_SUMMARY_INFO WHERE p_date = '%s' AND p_hour = %s """ session.sql(delete_query % (earlier_date, earlier_hour)).collect() else: slack_alert_query = """ select POST_TO_SLACK('#data-alerts', 'ML_SONG_SUMMARY_INFO_HOURLY_INSERT: failed execution check snowflake task') """ session.sql(slack_alert_query).collect() raise Exception("failed execution")