#! /usr/bin/env python3 from email.mime.multipart import MIMEMultipart from email.mime.text import MIMEText from contextlib import closing from sql import sql_queries import snowflake.connector import psycopg2 import smtplib import argparse import sys import json import datetime with open('config.json', 'r') as f: config = json.load(f) sf_user = config['snowflake']["user"] sf_password = config['snowflake']["password"] sf_account = config['snowflake']["account"] sf_warehouse = config['snowflake']['warehouse'] db_name_main = config['postgres_main']["database"] db_user_main = config['postgres_main']["user"] db_password_main = config['postgres_main']["password"] db_host_main = config['postgres_main']["host"] db_port_main = config['postgres_main']["port"] db_name_slz = config['postgres_slz']["database"] db_user_slz = config['postgres_slz']["user"] db_password_slz = config['postgres_slz']["password"] db_host_slz = config['postgres_slz']["host"] db_port_slz = config['postgres_slz']["port"] local_db_name = config['failure_counter']["database"] local_db_user = config['failure_counter']["user"] local_db_password = config['failure_counter']["password"] local_db_host = config['failure_counter']["host"] local_db_port = config['failure_counter']["port"] failure_count_increase = "UPDATE failed_count SET failed_count = failed_count + 1, last = 'FAILED', last_updated = NOW() WHERE count_id = 'aplch';" failure_count_to_zero = "UPDATE failed_count SET failed_count = 0, last = 'OK', last_updated = NOW() WHERE count_id = 'aplch';" def postgres_update_failure_count(sql_query): with closing(psycopg2.connect(database=local_db_name, user=local_db_user, password=local_db_password, host=local_db_host, port=local_db_port)) as connection: with connection.cursor() as cursor: cursor.execute(sql_query) connection.commit() def create_arg_parser(): parser = argparse.ArgumentParser() parser.add_argument('-m', '--mail', nargs='+', action='store', dest='recipients', type=str) parser.add_argument('-p', '--print', action="store_const", const=True) return parser def send_mail(recipients: str, message: str, status: str) -> None: # Create a multipart message msg = MIMEMultipart() body_part = MIMEText(message, 'plain') msg['Subject'] = f"[{status}][PROD] Apple Charts" msg['From'] = "support.sme@dataart.com" msg['To'] = recipients # Add body to email msg.attach(body_part) # Create SMTP object server = smtplib.SMTP('relay1.dataart.com', 25) server.sendmail(msg['From'], msg['To'].split(','), msg.as_string()) server.quit() def snowflake_fetch_data(sql_query) -> dict: ctx = snowflake.connector.connect( user=sf_user, password=sf_password, account=sf_account, warehouse=sf_warehouse ) cs = ctx.cursor(snowflake.connector.DictCursor) try: cs.execute(sql_query) result_dict = cs.fetchall() finally: cs.close() ctx.close() return result_dict def postgres_fetch_data_main(sql_query) -> list: result_list = [] with closing(psycopg2.connect(database=db_name_main, user=db_user_main, password=db_password_main, host=db_host_main, port=db_port_main)) as connection: with connection.cursor() as cursor: cursor.execute(sql_query) columns = cursor.description rows = cursor.fetchall() for row in rows: tmp = {} for i in range(len(columns)): tmp[columns[i][0]] = row[i] result_list.append(tmp) return result_list def postgres_fetch_data_slz(sql_query) -> list: result_list = [] with closing(psycopg2.connect(database=db_name_slz, user=db_user_slz, password=db_password_slz, host=db_host_slz, port=db_port_slz)) as connection: with connection.cursor() as cursor: cursor.execute(sql_query) columns = cursor.description rows = cursor.fetchall() for row in rows: tmp = {} for i in range(len(columns)): tmp[columns[i][0]] = row[i] result_list.append(tmp) return result_list def check_1() -> dict: result_dict = {} full_log = {} ### Check 1 SLZ status (expected result: 0 records)################################################ tmp_pg_slz_status_check_1 = postgres_fetch_data_slz(sql_queries.pg_slz_status_check_1) print(f"SLZ_STATUS_CHECK_1: - {tmp_pg_slz_status_check_1}") #Check if tmp_pg_slz_status_check_1: full_log["SLZ_STATUS_CHECK_1"] = {} result_dict["SLZ_STATUS_CHECK_1"] = {} for item in tmp_pg_slz_status_check_1: full_log["SLZ_STATUS_CHECK_1"].update({item['REPORT_DATE'].strftime("%Y-%m-%d"): item['MISSING_REPORTS'].lower()}) print(f"{item['REPORT_DATE']} - {item['MISSING_REPORTS']}") result_dict["SLZ_STATUS_CHECK_1"].update({item['REPORT_DATE'].strftime("%Y-%m-%d"): item['MISSING_REPORTS'].lower()}) else: full_log["SLZ_STATUS_CHECK_1"] = {'SLZ_STATUS_CHECK_1': 'Query returned 0 results'} return result_dict, full_log def check_2_to_22() -> dict: result_dict = {} full_log = {} ### ### Check 2 SLZ → RAW (expected result: is bigger then 0 )################################################ ### tmp_sf_slz_raw_check_2 = snowflake_fetch_data(sql_queries.sf_slz_raw_check_2) ### sf_slz_raw_check_2_count = int(tmp_sf_slz_raw_check_2[0]['count']) ### print(f"SLZ_TO_RAW_CHECK_2: - count: {sf_slz_raw_check_2_count}") ### full_log["SLZ_TO_RAW_CHECK_2"] = {} ### ### #Check ### if sf_slz_raw_check_2_count == 0: ### result_dict["SLZ_TO_RAW_CHECK_2"] = {} ### result_dict["SLZ_TO_RAW_CHECK_2"] = {'SLZ_TO_RAW_CHECK_2': 'FAILED: Query returned 0'} ### full_log["SLZ_TO_RAW_CHECK_2"] = result_dict["SLZ_TO_RAW_CHECK_2"] ### else: ### full_log["SLZ_TO_RAW_CHECK_2"] = {'SLZ_TO_RAW_CHECK_2': 'OK: Query returned more than 0'} ### ### ### Check 3 SLZ → RAW (expected result: 0 records) ################################################ ### tmp_sf_slz_raw_check_3 = snowflake_fetch_data(sql_queries.sf_slz_raw_check_3) ### sf_slz_raw_check_3_count = int(tmp_sf_slz_raw_check_3[0]['count']) ### print(f"SLZ_TO_RAW_CHECK_3: - count:{sf_slz_raw_check_3_count}") ### full_log["SLZ_TO_RAW_CHECK_3"] = {} ### ### #Check ### if sf_slz_raw_check_3_count > 0: ### result_dict["SLZ_TO_RAW_CHECK_3"] = {} ### result_dict["SLZ_TO_RAW_CHECK_3"] = {'SLZ_TO_RAW_CHECK_3': 'FAILED: Query returned more than 0'} ### full_log["SLZ_TO_RAW_CHECK_3"] = result_dict["SLZ_TO_RAW_CHECK_3"] ### else: ### full_log["SLZ_TO_RAW_CHECK_3"] = {'SLZ_TO_RAW_CHECK_3': 'OK: Query returned 0'} ### Check 4 RAW → MAIN - Data meta table status - Import status successful (expected result: Counts match inside every step) ################################################ tmp_sf_raw_main_check_4 = snowflake_fetch_data(sql_queries.sf_raw_main_check_4) sf_raw_main_check_4_src = int(tmp_sf_raw_main_check_4[0]['count']) sf_raw_main_check_4_dest = int(tmp_sf_raw_main_check_4[1]['count']) print(f"SRC_SF_RAW_CHECK_4: - {sf_raw_main_check_4_src}") print(f"DEST_SF_MAIN_CHECK_4: - {sf_raw_main_check_4_dest}") full_log["RAW_TO_MAIN_CHECK_4"] = {'SRC_SF_RAW': sf_raw_main_check_4_src, 'DEST_SF_MAIN': sf_raw_main_check_4_dest} if sf_raw_main_check_4_src != sf_raw_main_check_4_dest: result_dict["RAW_TO_MAIN_CHECK_4"] = full_log["RAW_TO_MAIN_CHECK_4"] ### Check 5 RAW → SYS - Data meta table status - Main sys table check (expected result: 0 records) ################################################ tmp_sf_raw_sys_check_5 = snowflake_fetch_data(sql_queries.sf_raw_sys_check_5) sf_raw_sys_check_5_count = int(tmp_sf_raw_sys_check_5[0]['count']) print(f"RAW_TO_SYS_CHECK_5: - count:{sf_raw_sys_check_5_count}") full_log["RAW_TO_SYS_CHECK_5"] = {} #Check if sf_raw_sys_check_5_count > 0: result_dict["RAW_TO_SYS_CHECK_5"] = {} result_dict["RAW_TO_SYS_CHECK_5"] = {'RAW_TO_SYS_CHECK_5': 'FAILED: Query returned more than 0'} full_log["RAW_TO_SYS_CHECK_5"] = result_dict["RAW_TO_SYS_CHECK_5"] else: full_log["RAW_TO_SYS_CHECK_5"] = {'RAW_TO_SYS_CHECK_5': 'OK: Query returned 0'} ### Check 6 RAW → MAIN - Data transfer status - Data consistency - DIM_APPLE_MUSIC_ARTIST (expected result: is bigger then 0 )################################################ tmp_sf_raw_main_check_6 = snowflake_fetch_data(sql_queries.sf_raw_main_check_6) sf_raw_main_check_6_src = int(tmp_sf_raw_main_check_6[0]['count']) sf_raw_main_check_6_dest = int(tmp_sf_raw_main_check_6[1]['count']) print(f"SRC_SF_RAW_CHECK_6: - {sf_raw_main_check_6_src}") print(f"DEST_SF_MAIN_CHECK_6: - {sf_raw_main_check_6_dest}") full_log["RAW_TO_MAIN_CHECK_6"] = {'SRC_SF_RAW': sf_raw_main_check_6_src, 'DEST_SF_MAIN': sf_raw_main_check_6_dest} if sf_raw_main_check_6_src != sf_raw_main_check_6_dest: result_dict["RAW_TO_MAIN_CHECK_6"] = full_log["RAW_TO_MAIN_CHECK_6"] ### Check 7 RAW → MAIN - Data transfer status - Data consistency - DIM_APPLE_MUSIC_CHART (expected result: Counts match inside every step) ################################################ tmp_sf_raw_main_check_7 = snowflake_fetch_data(sql_queries.sf_raw_main_check_7) sf_raw_main_check_7_src = int(tmp_sf_raw_main_check_7[0]['count']) sf_raw_main_check_7_dest = int(tmp_sf_raw_main_check_7[1]['count']) print(f"SRC_SF_RAW_CHECK_7: - {sf_raw_main_check_7_src}") print(f"DEST_SF_MAIN_CHECK_7: - {sf_raw_main_check_7_dest}") full_log["RAW_TO_MAIN_CHECK_7"] = {'SRC_SF_RAW': sf_raw_main_check_7_src, 'DEST_SF_MAIN': sf_raw_main_check_7_dest} if sf_raw_main_check_7_src != sf_raw_main_check_7_dest: result_dict["RAW_TO_MAIN_CHECK_7"] = full_log["RAW_TO_MAIN_CHECK_7"] ### Check 8 RAW → MAIN - Data transfer status - Data consistency - DIM_APPLE_MUSIC_TRACK (expected result: Counts match inside every step) ################################################ tmp_sf_raw_main_check_8 = snowflake_fetch_data(sql_queries.sf_raw_main_check_8) sf_raw_main_check_8_src = int(tmp_sf_raw_main_check_8[0]['count']) sf_raw_main_check_8_dest = int(tmp_sf_raw_main_check_8[1]['count']) print(f"SRC_SF_RAW_CHECK_8: - {sf_raw_main_check_8_src}") print(f"DEST_SF_MAIN_CHECK_8: - {sf_raw_main_check_8_dest}") full_log["RAW_TO_MAIN_CHECK_8"] = {'SRC_SF_RAW': sf_raw_main_check_8_src, 'DEST_SF_MAIN': sf_raw_main_check_8_dest} if sf_raw_main_check_8_src != sf_raw_main_check_8_dest: result_dict["RAW_TO_MAIN_CHECK_8"] = full_log["RAW_TO_MAIN_CHECK_8"] ### Check 9 RAW → MAIN - Data transfer status - Data consistency - DIM_APPLE_MUSIC_TRACK_LOCALIZED (expected result: Counts match inside every step) ################################################ tmp_sf_raw_main_check_9 = snowflake_fetch_data(sql_queries.sf_raw_main_check_9) sf_raw_main_check_9_src = int(tmp_sf_raw_main_check_9[0]['count']) sf_raw_main_check_9_dest = int(tmp_sf_raw_main_check_9[1]['count']) print(f"SRC_SF_RAW_CHECK_9: - {sf_raw_main_check_9_src}") print(f"DEST_SF_MAIN_CHECK_9: - {sf_raw_main_check_9_dest}") full_log["RAW_TO_MAIN_CHECK_9"] = {'SRC_SF_RAW': sf_raw_main_check_9_src, 'DEST_SF_MAIN': sf_raw_main_check_9_dest} if sf_raw_main_check_9_src != sf_raw_main_check_9_dest: result_dict["RAW_TO_MAIN_CHECK_9"] = full_log["RAW_TO_MAIN_CHECK_9"] ### Check 10 RAW→ MAIN - Data transfer status - Data consistency - DIM_APPLE_MUSIC_TRACK_PREVIEW (expected result: Counts match inside every step) ################################################ tmp_sf_raw_main_check_10 = snowflake_fetch_data(sql_queries.sf_raw_main_check_10) sf_raw_main_check_10_src = int(tmp_sf_raw_main_check_10[0]['count']) sf_raw_main_check_10_dest = int(tmp_sf_raw_main_check_10[1]['count']) print(f"SRC_SF_RAW_CHECK_10: - {sf_raw_main_check_10_src}") print(f"DEST_SF_MAIN_CHECK_10: - {sf_raw_main_check_10_dest}") full_log["RAW_TO_MAIN_CHECK_10"] = {'SRC_SF_RAW': sf_raw_main_check_10_src, 'DEST_SF_MAIN': sf_raw_main_check_10_dest} if sf_raw_main_check_10_src != sf_raw_main_check_10_dest: result_dict["RAW_TO_MAIN_CHECK_10"] = full_log["RAW_TO_MAIN_CHECK_10"] ### Check 11 RAW→ MAIN - Data transfer status - Data consistency - LINK_APPLE_MUSIC_TRACK_ARTIST (expected result: Counts match inside every step) ################################################ tmp_sf_raw_main_check_11 = snowflake_fetch_data(sql_queries.sf_raw_main_check_11) sf_raw_main_check_11_src = int(tmp_sf_raw_main_check_11[0]['count']) sf_raw_main_check_11_dest = int(tmp_sf_raw_main_check_11[1]['count']) print(f"SRC_SF_RAW_CHECK_11: - {sf_raw_main_check_11_src}") print(f"DEST_SF_MAIN_CHECK_11: - {sf_raw_main_check_11_dest}") full_log["RAW_TO_MAIN_CHECK_11"] = {'SRC_SF_RAW': sf_raw_main_check_11_src, 'DEST_SF_MAIN': sf_raw_main_check_11_dest} if sf_raw_main_check_11_src != sf_raw_main_check_11_dest: result_dict["RAW_TO_MAIN_CHECK_11"] = full_log["RAW_TO_MAIN_CHECK_11"] ### Check 12 RAW → MAIN - Data transfer status - Data consistency - FACT_APPLE_MUSIC_CHART_TRACK_PUBLIC (expected result: Counts match inside every step) ################################################ tmp_sf_raw_main_check_12 = snowflake_fetch_data(sql_queries.sf_raw_main_check_12) sf_raw_main_check_12_src = int(tmp_sf_raw_main_check_12[0]['count']) sf_raw_main_check_12_dest = int(tmp_sf_raw_main_check_12[1]['count']) print(f"SRC_SF_RAW_CHECK_12: - {sf_raw_main_check_12_src}") print(f"DEST_SF_MAIN_CHECK_12: - {sf_raw_main_check_12_dest}") full_log["RAW_TO_MAIN_CHECK_12"] = {'SRC_SF_RAW': sf_raw_main_check_12_src, 'DEST_SF_MAIN': sf_raw_main_check_12_dest} if sf_raw_main_check_12_src != sf_raw_main_check_12_dest: result_dict["RAW_TO_MAIN_CHECK_12"] = full_log["RAW_TO_MAIN_CHECK_12"] ### Check 13 MAIN → MAIN - Data transfer status - Data consistency - FACT_APPLE_MUSIC_CHART_TRACK_PUBLIC (expected result: Counts match inside every step) ################################################ tmp_sf_main_main_check_13 = snowflake_fetch_data(sql_queries.sf_main_main_check_13) sf_main_main_check_13_src = int(tmp_sf_main_main_check_13[0]['count']) sf_main_main_check_13_dest = int(tmp_sf_main_main_check_13[1]['count']) print(f"SRC_SF_MAIN_CHECK_13: - {sf_main_main_check_13_src}") print(f"DEST_SF_MAIN_CHECK_13: - {sf_main_main_check_13_dest}") full_log["MAIN_TO_MAIN_CHECK_13"] = {'SRC_SF_MAIN': sf_main_main_check_13_src, 'DEST_SF_MAIN': sf_main_main_check_13_dest} if sf_main_main_check_13_src != sf_main_main_check_13_dest: result_dict["MAIN_TO_MAIN_CHECK_13"] = full_log["MAIN_TO_MAIN_CHECK_13"] ### Check 14 MAIN → MAIN - Data transfer status - Data consistency - FACT_APPLE_MUSIC_CHART_TRACK_LIFETIME (expected result: Counts match inside every step) ################################################ tmp_sf_main_main_check_14 = snowflake_fetch_data(sql_queries.sf_main_main_check_14) sf_main_main_check_14_src = int(tmp_sf_main_main_check_14[0]['count']) sf_main_main_check_14_dest = int(tmp_sf_main_main_check_14[1]['count']) print(f"SRC_SF_MAIN_CHECK_14: - {sf_main_main_check_14_src}") print(f"DEST_SF_MAIN_CHECK_14: - {sf_main_main_check_14_dest}") full_log["MAIN_TO_MAIN_CHECK_14"] = {'SRC_SF_MAIN': sf_main_main_check_14_src, 'DEST_SF_MAIN': sf_main_main_check_14_dest} if sf_main_main_check_14_src != sf_main_main_check_14_dest: result_dict["MAIN_TO_MAIN_CHECK_14"] = full_log["MAIN_TO_MAIN_CHECK_14"] ### Check 15 MAIN → ETL - Data transfer status - Data consistency - DIM_APPLE_MUSIC_ARTIST (expected result: Counts match inside every step) ################################################ tmp_sf_main_etl_check_15 = snowflake_fetch_data(sql_queries.sf_main_etl_check_15) sf_main_etl_check_15_src = int(tmp_sf_main_etl_check_15[0]['count']) sf_main_etl_check_15_dest = int(tmp_sf_main_etl_check_15[1]['count']) print(f"SRC_SF_MAIN_CHECK_15: - {sf_main_etl_check_15_src}") print(f"DEST_SF_ETL_CHECK_15: - {sf_main_etl_check_15_dest}") full_log["MAIN_TO_ETL_CHECK_15"] = {'SRC_SF_MAIN': sf_main_etl_check_15_src, 'DEST_SF_ETL': sf_main_etl_check_15_dest} if sf_main_etl_check_15_src != sf_main_etl_check_15_dest: result_dict["MAIN_TO_ETL_CHECK_15"] = full_log["MAIN_TO_ETL_CHECK_15"] ### Check 16 MAIN → ETL - Data transfer status - Data consistency - DIM_APPLE_MUSIC_CHART (expected result: Counts match inside every step) ################################################ tmp_sf_main_etl_check_16 = snowflake_fetch_data(sql_queries.sf_main_etl_check_16) sf_main_etl_check_16_src = int(tmp_sf_main_etl_check_16[0]['count']) sf_main_etl_check_16_dest = int(tmp_sf_main_etl_check_16[1]['count']) print(f"SRC_SF_MAIN_CHECK_16: - {sf_main_etl_check_16_src}") print(f"DEST_SF_ETL_CHECK_16: - {sf_main_etl_check_16_dest}") full_log["MAIN_TO_ETL_CHECK_16"] = {'SRC_SF_MAIN': sf_main_etl_check_16_src, 'DEST_SF_ETL': sf_main_etl_check_16_dest} if sf_main_etl_check_16_src != sf_main_etl_check_16_dest: result_dict["MAIN_TO_ETL_CHECK_16"] = full_log["MAIN_TO_ETL_CHECK_16"] ### Check 17 MAIN → ETL - Data transfer status - Data consistency - DIM_APPLE_MUSIC_TRACK (expected result: Counts match inside every step) ################################################ tmp_sf_mian_etl_check_17 = snowflake_fetch_data(sql_queries.sf_mian_etl_check_17) sf_mian_etl_check_17_src = int(tmp_sf_mian_etl_check_17[0]['count']) sf_mian_etl_check_17_dest = int(tmp_sf_mian_etl_check_17[1]['count']) print(f"SRC_SF_MAIN_CHECK_17: - {sf_mian_etl_check_17_src}") print(f"DEST_SF_ETL_CHECK_17: - {sf_mian_etl_check_17_dest}") full_log["MAIN_TO_ETL_CHECK_17"] = {'SRC_SF_MAIN': sf_mian_etl_check_17_src, 'DEST_SF_ETL': sf_mian_etl_check_17_dest} if sf_mian_etl_check_17_src != sf_mian_etl_check_17_dest: result_dict["MAIN_TO_ETL_CHECK_17"] = full_log["MAIN_TO_ETL_CHECK_17"] ### Check 18 MAIN → ETL - Data transfer status - Data consistency - DIM_APPLE_MUSIC_TRACK_LOCALIZED (expected result: Counts match inside every step) ################################################ tmp_sf_mian_etl_check_18 = snowflake_fetch_data(sql_queries.sf_mian_etl_check_18) sf_mian_etl_check_18_src = int(tmp_sf_mian_etl_check_18[0]['count']) sf_mian_etl_check_18_dest = int(tmp_sf_mian_etl_check_18[1]['count']) print(f"SRC_SF_MAIN_CHECK_18: - {sf_mian_etl_check_18_src}") print(f"DEST_SF_ETL_CHECK_18: - {sf_mian_etl_check_18_dest}") full_log["MAIN_TO_ETL_CHECK_18"] = {'SRC_SF_MAIN': sf_mian_etl_check_18_src, 'DEST_SF_ETL': sf_mian_etl_check_18_dest} if sf_mian_etl_check_18_src != sf_mian_etl_check_18_dest: result_dict["MAIN_TO_ETL_CHECK_18"] = full_log["MAIN_TO_ETL_CHECK_18"] ### Check 19 MAIN → ETL - Data transfer status - Data consistency - DIM_APPLE_MUSIC_TRACK_PREVIEW (expected result: Counts match inside every step) ################################################ tmp_sf_mian_etl_check_19 = snowflake_fetch_data(sql_queries.sf_mian_etl_check_19) sf_mian_etl_check_19_src = int(tmp_sf_mian_etl_check_19[0]['count']) sf_mian_etl_check_19_dest = int(tmp_sf_mian_etl_check_19[1]['count']) print(f"SRC_SF_MAIN_CHECK_19: - {sf_mian_etl_check_19_src}") print(f"DEST_SF_ETL_CHECK_19: - {sf_mian_etl_check_19_dest}") full_log["MAIN_TO_ETL_CHECK_19"] = {'SRC_SF_MAIN': sf_mian_etl_check_19_src, 'DEST_SF_ETL': sf_mian_etl_check_19_dest} if sf_mian_etl_check_19_src != sf_mian_etl_check_19_dest: result_dict["MAIN_TO_ETL_CHECK_19"] = full_log["MAIN_TO_ETL_CHECK_19"] ### Check 20 MAIN → ETL - Data transfer status - Data consistency - FACT_APPLE_MUSIC_CHART_TRACK (expected result: Counts match inside every step) ################################################ tmp_sf_mian_etl_check_20 = snowflake_fetch_data(sql_queries.sf_mian_etl_check_20) sf_mian_etl_check_20_src = int(tmp_sf_mian_etl_check_20[0]['count']) sf_mian_etl_check_20_dest = int(tmp_sf_mian_etl_check_20[1]['count']) print(f"SRC_SF_MAIN_CHECK_20: - {sf_mian_etl_check_20_src}") print(f"DEST_SF_ETL_CHECK_20: - {sf_mian_etl_check_20_dest}") full_log["MAIN_TO_ETL_CHECK_20"] = {'SRC_SF_MAIN': sf_mian_etl_check_20_src, 'DEST_SF_ETL': sf_mian_etl_check_20_dest} if sf_mian_etl_check_20_src != sf_mian_etl_check_20_dest: result_dict["MAIN_TO_ETL_CHECK_20"] = full_log["MAIN_TO_ETL_CHECK_20"] ### Check 21 MAIN → ETL - Data transfer status - Data consistency - FACT_APPLE_MUSIC_CHART_TRACK_LIFETIME (expected result: Counts match inside every step) ################################################ tmp_sf_mian_etl_check_21 = snowflake_fetch_data(sql_queries.sf_mian_etl_check_21) sf_mian_etl_check_21_src = int(tmp_sf_mian_etl_check_21[0]['count']) sf_mian_etl_check_21_dest = int(tmp_sf_mian_etl_check_21[1]['count']) print(f"SRC_SF_MAIN_CHECK_21: - {sf_mian_etl_check_21_src}") print(f"DEST_SF_ETL_CHECK_21: - {sf_mian_etl_check_21_dest}") full_log["MAIN_TO_ETL_CHECK_21"] = {'SRC_SF_MAIN': sf_mian_etl_check_21_src, 'DEST_SF_ETL': sf_mian_etl_check_21_dest} if sf_mian_etl_check_21_src != sf_mian_etl_check_21_dest: result_dict["MAIN_TO_ETL_CHECK_21"] = full_log["MAIN_TO_ETL_CHECK_21"] ### Check 22 MAIN → ETL - Data transfer status - Data consistency - FACT_APPLE_MUSIC_CHART_TRACK_LIFETIME (expected result: Counts match inside every step) ################################################ tmp_sf_mian_etl_check_22 = snowflake_fetch_data(sql_queries.sf_mian_etl_check_22) sf_mian_etl_check_22_src = int(tmp_sf_mian_etl_check_22[0]['count']) sf_mian_etl_check_22_dest = int(tmp_sf_mian_etl_check_22[1]['count']) print(f"SRC_SF_MAIN_CHECK_22: - {sf_mian_etl_check_22_src}") print(f"DEST_SF_ETL_CHECK_22: - {sf_mian_etl_check_22_dest}") full_log["MAIN_TO_ETL_CHECK_22"] = {'SRC_SF_MAIN': sf_mian_etl_check_22_src, 'DEST_SF_ETL': sf_mian_etl_check_22_dest} if sf_mian_etl_check_22_src != sf_mian_etl_check_22_dest: result_dict["MAIN_TO_ETL_CHECK_22"] = full_log["MAIN_TO_ETL_CHECK_22"] return result_dict, full_log def check_23_to_31() -> dict: result_dict = {} full_log = {} ### Check 24 Data transfer status - Data consistency - dim_apple_music_artist_localized (expected result: Counts match between steps) ################################################ tmp_sf_dts_check_24 = snowflake_fetch_data(sql_queries.sf_dts_check_24) tmp_pg_mian_dts_check_24 = postgres_fetch_data_main(sql_queries.pg_mian_dts_check_24) sf_dts_check_24_count = int(tmp_sf_dts_check_24[0]['count']) pg_mian_dts_check_24_count = int(tmp_pg_mian_dts_check_24[0]['count']) print(f"SRC_SF_ETL_CHECK_24: - {sf_dts_check_24_count}") print(f"DEST_PG_CHECK_24: - {pg_mian_dts_check_24_count}") full_log["DTS_DC_CHECK_24"] = {'SRC_SF_ETL': sf_dts_check_24_count, 'DEST_PG': pg_mian_dts_check_24_count} if sf_dts_check_24_count != pg_mian_dts_check_24_count: result_dict["DTS_DC_CHECK_24"] = full_log["DTS_DC_CHECK_24"] ### Check 25 Data transfer status - Data consistency - dim_apple_music_track_localized (expected result: Counts match between steps) ################################################ tmp_sf_dts_check_25 = snowflake_fetch_data(sql_queries.sf_dts_check_25) tmp_pg_mian_dts_check_25 = postgres_fetch_data_main(sql_queries.pg_mian_dts_check_25) sf_dts_check_25_count = int(tmp_sf_dts_check_25[0]['count']) pg_mian_dts_check_25_count = int(tmp_pg_mian_dts_check_25[0]['count']) print(f"SRC_SF_ETL_CHECK_25: - {sf_dts_check_25_count}") print(f"DEST_PG_CHECK_25: - {pg_mian_dts_check_25_count}") full_log["DTS_DC_CHECK_25"] = {'SRC_SF_ETL': sf_dts_check_25_count, 'DEST_PG': pg_mian_dts_check_25_count} if sf_dts_check_25_count != pg_mian_dts_check_25_count: result_dict["DTS_DC_CHECK_25"] = full_log["DTS_DC_CHECK_25"] ### Check 26 Data transfer status - Data consistency - fact_apple_music_chart_track (expected result: Counts match between steps) ################################################ tmp_sf_dts_check_26 = snowflake_fetch_data(sql_queries.sf_dts_check_26) tmp_pg_mian_dts_check_26 = postgres_fetch_data_main(sql_queries.pg_mian_dts_check_26) sf_dts_check_26_count = int(tmp_sf_dts_check_26[0]['count']) pg_mian_dts_check_26_count = int(tmp_pg_mian_dts_check_26[0]['count']) print(f"SRC_SF_ETL_CHECK_26: - {sf_dts_check_26_count}") print(f"DEST_PG_CHECK_26: - {pg_mian_dts_check_26_count}") full_log["DTS_DC_CHECK_26"] = {'SRC_SF_ETL': sf_dts_check_26_count, 'DEST_PG': pg_mian_dts_check_26_count} if sf_dts_check_26_count != pg_mian_dts_check_26_count: result_dict["DTS_DC_CHECK_26"] = full_log["DTS_DC_CHECK_26"] ### Check 27 Data transfer status - Data consistency - fact_apple_music_chart_track_lifetime (expected result: Counts match between steps) ################################################ tmp_sf_dts_check_27 = snowflake_fetch_data(sql_queries.sf_dts_check_27) tmp_pg_mian_dts_check_27 = postgres_fetch_data_main(sql_queries.pg_mian_dts_check_27) sf_dts_check_27_count = int(tmp_sf_dts_check_27[0]['count']) pg_mian_dts_check_27_count = int(tmp_pg_mian_dts_check_27[0]['count']) print(f"SRC_SF_ETL_CHECK_27: - {sf_dts_check_27_count}") print(f"DEST_PG_CHECK_27: - {pg_mian_dts_check_27_count}") full_log["DTS_DC_CHECK_27"] = {'SRC_SF_ETL': sf_dts_check_27_count, 'DEST_PG': pg_mian_dts_check_27_count} if sf_dts_check_27_count != pg_mian_dts_check_27_count: result_dict["DTS_DC_CHECK_27"] = full_log["DTS_DC_CHECK_27"] ### Check 28 Data transfer status - Data consistency - dim_apple_music_chart (expected result: Counts match between steps) ################################################ tmp_sf_dts_check_28 = snowflake_fetch_data(sql_queries.sf_dts_check_28) tmp_pg_mian_dts_check_28 = postgres_fetch_data_main(sql_queries.pg_mian_dts_check_28) sf_dts_check_28_count = int(tmp_sf_dts_check_28[0]['count']) pg_mian_dts_check_28_count = int(tmp_pg_mian_dts_check_28[0]['count']) print(f"SRC_SF_ETL_CHECK_28: - {sf_dts_check_28_count}") print(f"DEST_PG_CHECK_28: - {pg_mian_dts_check_28_count}") full_log["DTS_DC_CHECK_28"] = {'SRC_SF_ETL': sf_dts_check_28_count, 'DEST_PG': pg_mian_dts_check_28_count} if sf_dts_check_28_count != pg_mian_dts_check_28_count: result_dict["DTS_DC_CHECK_28"] = full_log["DTS_DC_CHECK_28"] ### Check 29 Data transfer status - Data consistency - link_apple_music_track_artist (expected result: Counts match between steps) ################################################ tmp_sf_dts_check_29 = snowflake_fetch_data(sql_queries.sf_dts_check_29) tmp_pg_mian_dts_check_29 = postgres_fetch_data_main(sql_queries.pg_mian_dts_check_29) sf_dts_check_29_count = int(tmp_sf_dts_check_29[0]['count']) pg_mian_dts_check_29_count = int(tmp_pg_mian_dts_check_29[0]['count']) print(f"SRC_SF_ETL_CHECK_29: - {sf_dts_check_29_count}") print(f"DEST_PG_CHECK_29: - {pg_mian_dts_check_29_count}") full_log["DTS_DC_CHECK_29"] = {'SRC_SF_ETL': sf_dts_check_29_count, 'DEST_PG': pg_mian_dts_check_29_count} if sf_dts_check_29_count != pg_mian_dts_check_29_count: result_dict["DTS_DC_CHECK_29"] = full_log["DTS_DC_CHECK_29"] ### Check 30 Data transfer status - Data consistency - dim_apple_music_artist (expected result: Counts match between steps) ################################################ tmp_sf_dts_check_30 = snowflake_fetch_data(sql_queries.sf_dts_check_30) tmp_pg_mian_dts_check_30 = postgres_fetch_data_main(sql_queries.pg_mian_dts_check_30) sf_dts_check_30_count = int(tmp_sf_dts_check_30[0]['count']) pg_mian_dts_check_30_count = int(tmp_pg_mian_dts_check_30[0]['count']) print(f"SRC_SF_ETL_CHECK_30: - {sf_dts_check_30_count}") print(f"DEST_PG_CHECK_30: - {pg_mian_dts_check_30_count}") full_log["DTS_DC_CHECK_30"] = {'SRC_SF_ETL': sf_dts_check_30_count, 'DEST_PG': pg_mian_dts_check_30_count} if sf_dts_check_30_count != pg_mian_dts_check_30_count: result_dict["DTS_DC_CHECK_30"] = full_log["DTS_DC_CHECK_30"] ### Check 31 Data transfer status - Data consistency - dim_apple_music_track (expected result: Counts match between steps) ################################################ tmp_sf_dts_check_31 = snowflake_fetch_data(sql_queries.sf_dts_check_31) tmp_pg_mian_dts_check_31 = postgres_fetch_data_main(sql_queries.pg_mian_dts_check_31) sf_dts_check_31_count = int(tmp_sf_dts_check_31[0]['count']) pg_mian_dts_check_31_count = int(tmp_pg_mian_dts_check_31[0]['count']) print(f"SRC_SF_ETL_CHECK_31: - {sf_dts_check_31_count}") print(f"DEST_PG_CHECK_31: - {pg_mian_dts_check_31_count}") full_log["DTS_DC_CHECK_31"] = {'SRC_SF_ETL': sf_dts_check_31_count, 'DEST_PG': pg_mian_dts_check_31_count} if sf_dts_check_31_count != pg_mian_dts_check_31_count: result_dict["DTS_DC_CHECK_31"] = full_log["DTS_DC_CHECK_31"] ### Check 32 check if all markets \ days exist (expected result: 0 records) ################################################ tmp_pg_markets_days_check_32 = postgres_fetch_data_main(sql_queries.pg_markets_days_check_32) pg_markets_days_check_32_count = int(tmp_pg_markets_days_check_32[0]['count']) print(f"MARKETS_DAYS_CHECK_32: - count:{pg_markets_days_check_32_count}") full_log["MARKETS_DAYS_CHECK_32"] = {} #Check if pg_markets_days_check_32_count != 0: result_dict["MARKETS_DAYS_CHECK_32"] = {} result_dict["MARKETS_DAYS_CHECK_32"] = {'MARKETS_DAYS_CHECK_32': 'FAILED: Query does not return 0'} full_log["MARKETS_DAYS_CHECK_32"] = result_dict["MARKETS_DAYS_CHECK_32"] else: full_log["MARKETS_DAYS_CHECK_32"] = {'MARKETS_DAYS_CHECK_32': 'OK: Query returned 0'} ### Check 33 check if there is no new markets (expected result: 0 records) ################################################ tmp_sf_new_markets_check_33 = snowflake_fetch_data(sql_queries.sf_new_markets_check_33) sf_new_markets_check_33_count = int(tmp_sf_new_markets_check_33[0]['count']) print(f"NEW_MARKETS_CHECK_33: - count:{sf_new_markets_check_33_count}") full_log["NEW_MARKETS_CHECK_33"] = {} #Check if sf_new_markets_check_33_count != 0: result_dict["NEW_MARKETS_CHECK_33"] = {} result_dict["NEW_MARKETS_CHECK_33"] = {'NEW_MARKETS_CHECK_33': 'FAILED: Query does not return 0'} full_log["NEW_MARKETS_CHECK_33"] = result_dict["NEW_MARKETS_CHECK_33"] else: full_log["NEW_MARKETS_CHECK_33"] = {'NEW_MARKETS_CHECK_33': 'OK: Query returned 0'} return result_dict, full_log ################# run the all steps ###################### def run_steps() -> dict: message_dict = {'name':"[PROD] Apple Charts",'status':"",'date':"", 'info':"", 'details':"", 'full_log':""} # set date according to sql searching date, currently it's `today` message_dict["date"] = (datetime.datetime.now() - datetime.timedelta(days=1)).strftime("%Y-%m-%d") result_dict_1, monitoring_1_log = check_1() result_dict_2_to_22, monitoring_2_to_22_log = check_2_to_22() result_dict_23_to_31, monitoring_23_to_31_log = check_23_to_31() message_dict['full_log'] = {**monitoring_1_log, **monitoring_2_to_22_log, **monitoring_23_to_31_log} message_dict['details'] = {**result_dict_1, **result_dict_2_to_22, **result_dict_23_to_31} if result_dict_2_to_22 or result_dict_1 or result_dict_23_to_31: postgres_update_failure_count(failure_count_increase) if result_dict_1: # Some UOWs aren’t completed message_dict['status'] = "FAILED" message_dict['info'] = "Some UOWs are not completed" message_dict['check'] = "1" elif result_dict_2_to_22: # Data is synced with Fivetran and partially available in API message_dict['status'] = "FAILED" message_dict['info'] = "Data is partially ingested to Snowflake" message_dict['check'] = "2" elif result_dict_23_to_31: message_dict['status'] = "FAILED" message_dict['info'] = "Data is partially available in API" message_dict['check'] = "3" else: postgres_update_failure_count(failure_count_to_zero) message_dict['status'] = "OK" message_dict['info'] = "Data is available in API" return message_dict if __name__ == '__main__': # # check command line parameters and print or email result parser = create_arg_parser() arg_list = parser.parse_args() if len(sys.argv[1:]) > 0: result = run_steps() with open('/tmp/Apple_daily_log.json', 'w') as json_file: json.dump(result, json_file, indent=4) if arg_list.recipients: recipients = arg_list.recipients[0] send_mail(recipients=recipients, message=json.dumps(result, indent=4), status=result['status']) if arg_list.print: print(json.dumps(result, indent=4)) else: parser.print_help()