#! /usr/bin/env python3 from email.mime.multipart import MIMEMultipart from email.mime.text import MIMEText from contextlib import closing 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 = config['postgres']["database"] db_user = config['postgres']["user"] db_password = config['postgres']["password"] db_host = config['postgres']["host"] db_port = config['postgres']["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"] monitoring_queries = { "DS_CHARTMETRICK_VS_RAW": { "DS_CHARTMETRICK_VS_RAW_STEP_1": "select count(*) as count,'raw' as tbl from DELPHI_CHARTMETRIC_RAW.CHARTS_SHAZAM.SHAZAM_CHART WHERE TIMESTP <= dateadd(dd, -1, current_date()) and TIMESTP != '2023-06-10' and not ((trim(COUNTRY) = 'UA'and timestp = '2021-07-27' and city_id is not null and city_id = '795874') or (trim(COUNTRY) = 'RU'and timestp = '2021-07-27' and city_id is not null and city_id = '5372') or (trim(COUNTRY) = 'BY'and timestp = '2021-09-07' and city_id is null) or (trim(COUNTRY) = 'TR'and timestp = '2021-09-14' and city_id is not null and city_id = '6101')) union select count(*),'chart' as src from DS_CHARTMETRIC.RAW_DATA.SHAZAM_CHART WHERE TIMESTP <= dateadd(dd, -1, current_date()) and TIMESTP != '2023-06-10';", "DS_CHARTMETRICK_VS_RAW_STEP_2": "select COUNT(*) as count,'raw' as tbl from DELPHI_CHARTMETRIC_RAW.CHARTS_SHAZAM.SHAZAM_STAT WHERE TIMESTP <= dateadd(dd, -1, current_date()) union select COUNT(*),'chart' from DS_CHARTMETRIC.RAW_DATA.SHAZAM_STAT WHERE TIMESTP <= dateadd(dd, -1, current_date());" }, "RAW_VS_MAIN": { "RAW_VS_MAIN_STEP_1": "select COUNT(*) as count,'raw' as tbl from DELPHI_CHARTMETRIC_RAW.CHARTS_SHAZAM.SHAZAM union select COUNT(*),'main' from DELPHI_CHARTMETRIC_MAIN.CHARTS_SHAZAM.DIM_SHAZAM_CHART_TRACK_META;", "RAW_VS_MAIN_STEP_2": "select COUNT(*) as count,'raw' as tbl from (select * from (select * from DELPHI_CHARTMETRIC_RAW.CHARTS_SHAZAM.SHAZAM_CHART where CITY_ID is not null and RANK <= 200 qualify row_number()over (partition by (trim(COUNTRY), CITY_ID, RANK, TIMESTP) order by MODIFIED_AT desc, ID desc) =1) qualify row_number() over (partition by (trim(COUNTRY), CITY_ID, SHAZAM_TRACK_ID, TIMESTP) order by RANK ) = 1 ) union select COUNT(*), 'main' from DELPHI_CHARTMETRIC_MAIN.CHARTS_SHAZAM.FACT_SHAZAM_CHART_CITY;", "RAW_VS_MAIN_STEP_3": "select COUNT(*) as count,'raw' as tbl from (select * from (select * from DELPHI_CHARTMETRIC_RAW.CHARTS_SHAZAM.SHAZAM_CHART where CITY_ID is null and RANK <= 200 qualify row_number()over (partition by (trim(COUNTRY), RANK, TIMESTP) order by MODIFIED_AT desc, ID desc) =1) qualify row_number() over (partition by (trim(COUNTRY), SHAZAM_TRACK_ID, TIMESTP) order by RANK ) = 1 )union select COUNT(*), 'main' from DELPHI_CHARTMETRIC_MAIN.CHARTS_SHAZAM.FACT_SHAZAM_CHART_COUNTRY;", "RAW_VS_MAIN_STEP_4": "select COUNT(*) as count,'raw' as tbl from DELPHI_CHARTMETRIC_RAW.CHARTS_SHAZAM.SHAZAM_STAT union select COUNT(*),'main' from DELPHI_CHARTMETRIC_MAIN.CHARTS_SHAZAM.FACT_TOTAL_SHAZAMS;", "RAW_VS_MAIN_STEP_5": "select COUNT(*) as count,'raw' as tbl from DELPHI_CHARTMETRIC_MAIN.CHARTS_SHAZAM.SHAZAM_CHART_TRACK_SUMMARY_REPORT union select count(distinct coalesce(CITY_ID || '', trim(COUNTRY)), SHAZAM_TRACK_ID),'main' from (select * from (select * from DELPHI_CHARTMETRIC_RAW.CHARTS_SHAZAM.SHAZAM_CHART where RANK <= 200 qualify row_number() over (partition by (coalesce(CITY_ID || '', trim(COUNTRY)), RANK, TIMESTP) order by MODIFIED_AT desc, ID desc) = 1)qualify row_number()over (partition by (coalesce(CITY_ID || '', trim(COUNTRY)), SHAZAM_TRACK_ID, TIMESTP) order by RANK ) = 1 );", "RAW_VS_MAIN_STEP_6": "select count(distinct coalesce(CITY_ID || '', trim(COUNTRY))) as count,'raw' as tbl from DELPHI_CHARTMETRIC_RAW.CHARTS_SHAZAM.SHAZAM_CHART union select COUNT(*), 'main_teritory' from DELPHI_CHARTMETRIC_MAIN.CHARTS_SHAZAM.DIM_SHAZAM_TERRITORY union select COUNT(*), 'main_chart' from DELPHI_CHARTMETRIC_MAIN.CHARTS_SHAZAM.DIM_SHAZAM_CHART_META;" }, "MAIN_VS_ETL": { "MAIN_VS_ETL_STEP_1": "select COUNT(*) as count,'etl' as tbl from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.DIM_SHAZAM_CHART_META union select COUNT(*),'main' from DELPHI_CHARTMETRIC_MAIN.CHARTS_SHAZAM.DIM_SHAZAM_CHART_META;", "MAIN_VS_ETL_STEP_2": "select COUNT(*) as count,'etl' as tbl from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.DIM_SHAZAM_CHART_TRACK_META union select COUNT(*),'main' from DELPHI_CHARTMETRIC_MAIN.CHARTS_SHAZAM.DIM_SHAZAM_CHART_TRACK_META;", "MAIN_VS_ETL_STEP_3": "select COUNT(*) as count,'etl' as tbl from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.DIM_SHAZAM_TERRITORY union select COUNT(*),'main' from DELPHI_CHARTMETRIC_MAIN.CHARTS_SHAZAM.DIM_SHAZAM_TERRITORY;", "MAIN_VS_ETL_STEP_4": "select COUNT(*) as count,'etl' as tbl from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.FACT_SHAZAM_CHART_CITY union select COUNT(*),'main' from DELPHI_CHARTMETRIC_MAIN.CHARTS_SHAZAM.FACT_SHAZAM_CHART_CITY;", "MAIN_VS_ETL_STEP_5": "select COUNT(*) as count,'etl' as tbl from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.FACT_SHAZAM_CHART_COUNTRY union select COUNT(*),'main' from DELPHI_CHARTMETRIC_MAIN.CHARTS_SHAZAM.FACT_SHAZAM_CHART_COUNTRY;", "MAIN_VS_ETL_STEP_6": "select COUNT(*) as count,'etl' as tbl from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.FACT_TOTAL_SHAZAMS union select COUNT(*),'main' from DELPHI_CHARTMETRIC_MAIN.CHARTS_SHAZAM.FACT_TOTAL_SHAZAMS;", "MAIN_VS_ETL_STEP_7": "select COUNT(*) as count,'etl' as tbl from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.SHAZAM_CHART_TRACK_SUMMARY_REPORT union select COUNT(*),'main' from DELPHI_CHARTMETRIC_MAIN.CHARTS_SHAZAM.SHAZAM_CHART_TRACK_SUMMARY_REPORT;" }, "DTS_DC_DIMENTIONAL_TABLE": { "DIM_SHAZAM_CHARTS_SF": "select COUNT(*) as count from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.DIM_SHAZAM_CHART_META;", "DIM_SHAZAM_CHARTS_PG": "select COUNT(*) as count from chartmetric.dim_shazam_chart;", "DIM_SHAZAM_CITY_SF": "select COUNT(*) as count from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.DIM_SHAZAM_TERRITORY where SHAZAM_CITY_ID is not null;", "DIM_SHAZAM_CITY_PG": "select COUNT(*) as count from chartmetric.dim_shazam_city;", "DIM_SHAZAM_TRACK_SF": "select COUNT(*) as count from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.DIM_SHAZAM_CHART_TRACK_META;", "DIM_SHAZAM_TRACK_PG": "select COUNT(*) as count from chartmetric.dim_shazam_track;", "SHAZAM_CHART_TRACK_SUMMARY_REPORT_SF": "select COUNT(*) as count from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.SHAZAM_CHART_TRACK_SUMMARY_REPORT where TRACK_ID in (select TRACK_ID from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.DIM_SHAZAM_CHART_TRACK_META);", "SHAZAM_CHART_TRACK_SUMMARY_REPORT_PG": "select COUNT(*) as count from public.shazam_chart_track_summary_report sctsr;", "DIM_SHAZAM_CHART_ARTIST_SF": "select COUNT(*) as count from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.DIM_SHAZAM_CHART_ARTIST_ISRC;", "DIM_SHAZAM_CHART_ARTIST_PG": "select COUNT(*) as count from chartmetric.dim_shazam_chart_artist_isrc;", "DIM_SHAZAM_CHART_ARTIST_ITUNES_SF": "select COUNT(*) as count from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.DIM_SHAZAM_CHART_ARTIST_ITUNES;", "DIM_SHAZAM_CHART_ARTIST_ITUNES_PG": "select COUNT(*) as count from chartmetric.dim_shazam_chart_artist_itunes;" }, "DTS_DC_FACT_TABLE": { "FACT_SHAZAM_CHART_SF": "select sum(count) as count from ( select count(*) as count from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.FACT_SHAZAM_CHART_CITY where TRACK_ID in (select TRACK_ID from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.DIM_SHAZAM_CHART_TRACK_META) union select count(*)as count from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.FACT_SHAZAM_CHART_COUNTRY where TRACK_ID in (select TRACK_ID from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.DIM_SHAZAM_CHART_TRACK_META));", "FACT_SHAZAM_CHART_PG": "select count(*) as count from chartmetric.fact_shazam_chart;", "FACT_SHAZAM_TRACK_TOTAL_SHAZAMS_SF": "select count(*) as count from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.FACT_TOTAL_SHAZAMS where SHAZAM_TRACK_ID in ( select TRACK_ID from DELPHI_CHARTMETRIC_ETL.CHARTS_SHAZAM.DIM_SHAZAM_CHART_TRACK_META);", "FACT_SHAZAM_TRACK_TOTAL_SHAZAMS_PG": "select count(*) as count from chartmetric.fact_shazam_track_total_shazams;" } } failure_count_increase = "UPDATE failed_count SET failed_count = failed_count + 1, last = 'FAILED', last_updated = NOW() WHERE count_id = 'shch';" failure_count_to_zero = "UPDATE failed_count SET failed_count = 0, last = 'OK', last_updated = NOW() WHERE count_id = 'shch';" 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] Shazam 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(sql_query) -> list: result_list = [] with closing(psycopg2.connect(database=db_name, user=db_user, password=db_password, host=db_host, port=db_port)) 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_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 monitoring_123() -> dict: result_dict = {} full_log = {} # Check 1 DS_CHARTMETRICK_VS_RAW tmp_fetched_ds_chartmetrick_vs_raw_step_1 = snowflake_fetch_data(monitoring_queries["DS_CHARTMETRICK_VS_RAW"]["DS_CHARTMETRICK_VS_RAW_STEP_1"]) ds_raw_step_1 = int(tmp_fetched_ds_chartmetrick_vs_raw_step_1[0]['COUNT']) ds_chart_step_1 = int(tmp_fetched_ds_chartmetrick_vs_raw_step_1[1]['COUNT']) tmp_fetched_ds_chartmetrick_vs_raw_step_2 = snowflake_fetch_data(monitoring_queries["DS_CHARTMETRICK_VS_RAW"]["DS_CHARTMETRICK_VS_RAW_STEP_2"]) ds_raw_step_2 = int(tmp_fetched_ds_chartmetrick_vs_raw_step_2[0]['COUNT']) ds_chart_step_2 = int(tmp_fetched_ds_chartmetrick_vs_raw_step_2[1]['COUNT']) print(f"Check 1 DS_CHARTMETRICK_VS_RAW STEP 1 - {tmp_fetched_ds_chartmetrick_vs_raw_step_1}") print(f"Check 1 DS_CHARTMETRICK_VS_RAW STEP 2 - {tmp_fetched_ds_chartmetrick_vs_raw_step_2}") full_log["DS_CHARTMETRICK_VS_RAW_STEP_1"] = {'RAW_STEP_1': ds_raw_step_1, 'CHARTMETRICK_STEP_1': ds_chart_step_1} full_log["DS_CHARTMETRICK_VS_RAW_STEP_2"] = {'RAW_STEP_2': ds_raw_step_2, 'CHARTMETRICK_STEP_2': ds_chart_step_2} if ds_raw_step_1 != ds_chart_step_1: result_dict["DS_CHARTMETRICK_VS_RAW_STEP_1"] = full_log["DS_CHARTMETRICK_VS_RAW_STEP_1"] if ds_raw_step_2 != ds_chart_step_2: result_dict["DS_CHARTMETRICK_VS_RAW_STEP_2"] = full_log["DS_CHARTMETRICK_VS_RAW_STEP_2"] # Check 2 RAW_VS_MAIN tmp_fetched_raw_vs_main_step_1 = snowflake_fetch_data(monitoring_queries["RAW_VS_MAIN"]["RAW_VS_MAIN_STEP_1"]) ds_raw_step_1 = int(tmp_fetched_raw_vs_main_step_1[0]['COUNT']) ds_main_step_1 = int(tmp_fetched_raw_vs_main_step_1[1]['COUNT']) tmp_fetched_raw_vs_main_step_2 = snowflake_fetch_data(monitoring_queries["RAW_VS_MAIN"]["RAW_VS_MAIN_STEP_2"]) ds_raw_step_2 = int(tmp_fetched_raw_vs_main_step_2[0]['COUNT']) ds_main_step_2 = int(tmp_fetched_raw_vs_main_step_2[1]['COUNT']) tmp_fetched_raw_vs_main_step_3 = snowflake_fetch_data(monitoring_queries["RAW_VS_MAIN"]["RAW_VS_MAIN_STEP_3"]) ds_raw_step_3 = int(tmp_fetched_raw_vs_main_step_3[0]['COUNT']) ds_main_step_3 = int(tmp_fetched_raw_vs_main_step_3[1]['COUNT']) tmp_fetched_raw_vs_main_step_4 = snowflake_fetch_data(monitoring_queries["RAW_VS_MAIN"]["RAW_VS_MAIN_STEP_4"]) ds_raw_step_4 = int(tmp_fetched_raw_vs_main_step_4[0]['COUNT']) ds_main_step_4 = int(tmp_fetched_raw_vs_main_step_4[1]['COUNT']) tmp_fetched_raw_vs_main_step_5 = snowflake_fetch_data(monitoring_queries["RAW_VS_MAIN"]["RAW_VS_MAIN_STEP_5"]) ds_raw_step_5 = int(tmp_fetched_raw_vs_main_step_5[0]['COUNT']) ds_main_step_5 = int(tmp_fetched_raw_vs_main_step_5[1]['COUNT']) tmp_fetched_raw_vs_main_step_6 = snowflake_fetch_data(monitoring_queries["RAW_VS_MAIN"]["RAW_VS_MAIN_STEP_6"]) ds_raw_step_6 = int(tmp_fetched_raw_vs_main_step_6[0]['COUNT']) ds_main_teritory_step_6 = int(tmp_fetched_raw_vs_main_step_6[1]['COUNT']) ds_main_chart_step_6 = int(tmp_fetched_raw_vs_main_step_6[2]['COUNT']) print(f"Check 2 RAW_VS_MAIN STEP 1 - {tmp_fetched_raw_vs_main_step_1}") print(f"Check 2 RAW_VS_MAIN STEP 2 - {tmp_fetched_raw_vs_main_step_2}") print(f"Check 2 RAW_VS_MAIN STEP 3 - {tmp_fetched_raw_vs_main_step_3}") print(f"Check 2 RAW_VS_MAIN STEP 4 - {tmp_fetched_raw_vs_main_step_4}") print(f"Check 2 RAW_VS_MAIN STEP 5 - {tmp_fetched_raw_vs_main_step_5}") print(f"Check 2 RAW_VS_MAIN STEP 6 - {tmp_fetched_raw_vs_main_step_6}") full_log["RAW_VS_MAIN_STEP_1"] = {'RAW step 1': ds_raw_step_1, 'MAIN step 1': ds_main_step_1} full_log["RAW_VS_MAIN_STEP_2"] = {'RAW step 2': ds_raw_step_2, 'MAIN step 2': ds_main_step_2} full_log["RAW_VS_MAIN_STEP_3"] = {'RAW step 3': ds_raw_step_3, 'MAIN step 3': ds_main_step_3} full_log["RAW_VS_MAIN_STEP_4"] = {'RAW step 4': ds_raw_step_4, 'MAIN step 4': ds_main_step_4} full_log["RAW_VS_MAIN_STEP_5"] = {'RAW step 5': ds_raw_step_5, 'MAIN step 5': ds_main_step_5} full_log["RAW_VS_MAIN_STEP_6"] = {'RAW step 6': ds_raw_step_6, 'MAIN_TERITORY step 6': ds_main_teritory_step_6, 'MAIN CHART step 6': ds_main_chart_step_6} if ds_raw_step_1 != ds_main_step_1: result_dict["RAW_VS_MAIN_STEP_1"] = full_log["RAW_VS_MAIN_STEP_1"] if ds_raw_step_2 != ds_main_step_2: result_dict["RAW_VS_MAIN_STEP_2"] = full_log["RAW_VS_MAIN_STEP_2"] if ds_raw_step_3 != ds_main_step_3: result_dict["RAW_VS_MAIN_STEP_3"] = full_log["RAW_VS_MAIN_STEP_3"] if ds_raw_step_4 != ds_main_step_4: result_dict["RAW_VS_MAIN_STEP_4"] = full_log["RAW_VS_MAIN_STEP_4"] if ds_raw_step_5 != ds_main_step_5: result_dict["RAW_VS_MAIN_STEP_5"] = full_log["RAW_VS_MAIN_STEP_5"] if ds_raw_step_6 != ds_main_teritory_step_6 or ds_raw_step_6 != ds_main_chart_step_6 or ds_main_teritory_step_6 != ds_main_chart_step_6: result_dict["RAW_VS_MAIN_STEP_6"] = full_log["RAW_VS_MAIN_STEP_6"] # Check 2 MAIN_VS_ETL tmp_fetched_main_vs_etl_step_1 = snowflake_fetch_data(monitoring_queries["MAIN_VS_ETL"]["MAIN_VS_ETL_STEP_1"]) ds_etl_step_1 = int(tmp_fetched_main_vs_etl_step_1[0]['COUNT']) ds_main_step_1 = int(tmp_fetched_main_vs_etl_step_1[1]['COUNT']) tmp_fetched_main_vs_etl_step_2 = snowflake_fetch_data(monitoring_queries["MAIN_VS_ETL"]["MAIN_VS_ETL_STEP_2"]) ds_etl_step_2 = int(tmp_fetched_main_vs_etl_step_2[0]['COUNT']) ds_main_step_2 = int(tmp_fetched_main_vs_etl_step_2[1]['COUNT']) tmp_fetched_main_vs_etl_step_3 = snowflake_fetch_data(monitoring_queries["MAIN_VS_ETL"]["MAIN_VS_ETL_STEP_3"]) ds_etl_step_3 = int(tmp_fetched_main_vs_etl_step_3[0]['COUNT']) ds_main_step_3 = int(tmp_fetched_main_vs_etl_step_3[1]['COUNT']) tmp_fetched_main_vs_etl_step_4 = snowflake_fetch_data(monitoring_queries["MAIN_VS_ETL"]["MAIN_VS_ETL_STEP_4"]) ds_etl_step_4 = int(tmp_fetched_main_vs_etl_step_4[0]['COUNT']) ds_main_step_4 = int(tmp_fetched_main_vs_etl_step_4[1]['COUNT']) tmp_fetched_main_vs_etl_step_5 = snowflake_fetch_data(monitoring_queries["MAIN_VS_ETL"]["MAIN_VS_ETL_STEP_5"]) ds_etl_step_5 = int(tmp_fetched_main_vs_etl_step_5[0]['COUNT']) ds_main_step_5 = int(tmp_fetched_main_vs_etl_step_5[1]['COUNT']) tmp_fetched_main_vs_etl_step_6 = snowflake_fetch_data(monitoring_queries["MAIN_VS_ETL"]["MAIN_VS_ETL_STEP_6"]) ds_etl_step_6 = int(tmp_fetched_main_vs_etl_step_6[0]['COUNT']) ds_main_step_6 = int(tmp_fetched_main_vs_etl_step_6[1]['COUNT']) tmp_fetched_main_vs_etl_step_7 = snowflake_fetch_data(monitoring_queries["MAIN_VS_ETL"]["MAIN_VS_ETL_STEP_7"]) ds_etl_step_7 = int(tmp_fetched_main_vs_etl_step_7[0]['COUNT']) ds_main_step_7 = int(tmp_fetched_main_vs_etl_step_7[1]['COUNT']) print(f"Check 2 MAIN_VS_ETL STEP 1 - {tmp_fetched_main_vs_etl_step_1}") print(f"Check 2 MAIN_VS_ETL STEP 2 - {tmp_fetched_main_vs_etl_step_2}") print(f"Check 2 MAIN_VS_ETL STEP 3 - {tmp_fetched_main_vs_etl_step_3}") print(f"Check 2 MAIN_VS_ETL STEP 4 - {tmp_fetched_main_vs_etl_step_4}") print(f"Check 2 MAIN_VS_ETL STEP 5 - {tmp_fetched_main_vs_etl_step_5}") print(f"Check 2 MAIN_VS_ETL STEP 6 - {tmp_fetched_main_vs_etl_step_6}") print(f"Check 2 MAIN_VS_ETL STEP 7 - {tmp_fetched_main_vs_etl_step_7}") full_log["MAIN_VS_ETL_STEP_1"] = {'MAIN_step_1': ds_main_step_1, 'ETL_step_1': ds_etl_step_1} full_log["MAIN_VS_ETL_STEP_2"] = {'MAIN_step_2': ds_main_step_2, 'ETL_step_2': ds_etl_step_2} full_log["MAIN_VS_ETL_STEP_3"] = {'MAIN_step_3': ds_main_step_3, 'ETL_step_3': ds_etl_step_3} full_log["MAIN_VS_ETL_STEP_4"] = {'MAIN_step_4': ds_main_step_4, 'ETL_step_4': ds_etl_step_4} full_log["MAIN_VS_ETL_STEP_5"] = {'MAIN_step_5': ds_main_step_5, 'ETL_step_5': ds_etl_step_5} full_log["MAIN_VS_ETL_STEP_6"] = {'MAIN_step_6': ds_main_step_6, 'ETL_step_6': ds_etl_step_6} full_log["MAIN_VS_ETL_STEP_7"] = {'MAIN_step_7': ds_main_step_7, 'ETL_step_7': ds_etl_step_7} if ds_etl_step_1 != ds_main_step_1: result_dict["MAIN_VS_ETL_STEP_1"] = full_log["MAIN_VS_ETL_STEP_1"] if ds_etl_step_2 != ds_main_step_2: result_dict["MAIN_VS_ETL_STEP_2"] = full_log["MAIN_VS_ETL_STEP_2"] if ds_etl_step_3 != ds_main_step_3: result_dict["MAIN_VS_ETL_STEP_3"] = full_log["MAIN_VS_ETL_STEP_3"] if ds_etl_step_4 != ds_main_step_4: result_dict["MAIN_VS_ETL_STEP_4"] = full_log["MAIN_VS_ETL_STEP_4"] if ds_etl_step_5 != ds_main_step_5: result_dict["MAIN_VS_ETL_STEP_5"] = full_log["MAIN_VS_ETL_STEP_5"] if ds_etl_step_6 != ds_main_step_6: result_dict["MAIN_VS_ETL_STEP_6"] = full_log["MAIN_VS_ETL_STEP_6"] if ds_etl_step_7 != ds_main_step_7: result_dict["MAIN_VS_ETL_STEP_7"] = full_log["MAIN_VS_ETL_STEP_7"] return result_dict, full_log def monitoring_56() -> dict: result_dict = {} full_log = {} # Check 5 DTS_DC_DIMENTIONAL_TABLE dim_shahzam_charts_sf = int(snowflake_fetch_data(monitoring_queries["DTS_DC_DIMENTIONAL_TABLE"]["DIM_SHAZAM_CHARTS_SF"])[0]['COUNT']) dim_shahzam_charts_pg = int(postgres_fetch_data(monitoring_queries["DTS_DC_DIMENTIONAL_TABLE"]["DIM_SHAZAM_CHARTS_PG"])[0]['count']) dim_shahzam_city_sf = int(snowflake_fetch_data(monitoring_queries["DTS_DC_DIMENTIONAL_TABLE"]["DIM_SHAZAM_CITY_SF"])[0]['COUNT']) dim_shahzam_city_pg = int(postgres_fetch_data(monitoring_queries["DTS_DC_DIMENTIONAL_TABLE"]["DIM_SHAZAM_CITY_PG"])[0]['count']) dim_shahzam_track_sf = int(snowflake_fetch_data(monitoring_queries["DTS_DC_DIMENTIONAL_TABLE"]["DIM_SHAZAM_TRACK_SF"])[0]['COUNT']) dim_shahzam_track_pg = int(postgres_fetch_data(monitoring_queries["DTS_DC_DIMENTIONAL_TABLE"]["DIM_SHAZAM_TRACK_PG"])[0]['count']) dim_shahzam_track_summary_report_sf = int(snowflake_fetch_data(monitoring_queries["DTS_DC_DIMENTIONAL_TABLE"]["SHAZAM_CHART_TRACK_SUMMARY_REPORT_SF"])[0]['COUNT']) dim_shahzam_track_summary_report_pg = int(postgres_fetch_data(monitoring_queries["DTS_DC_DIMENTIONAL_TABLE"]["SHAZAM_CHART_TRACK_SUMMARY_REPORT_PG"])[0]['count']) dim_shahzam_charts_artist_sf = int(snowflake_fetch_data(monitoring_queries["DTS_DC_DIMENTIONAL_TABLE"]["DIM_SHAZAM_CHART_ARTIST_SF"])[0]['COUNT']) dim_shahzam_charts_artist_pg = int(postgres_fetch_data(monitoring_queries["DTS_DC_DIMENTIONAL_TABLE"]["DIM_SHAZAM_CHART_ARTIST_PG"])[0]['count']) dim_shahzam_charts_artist_itunes_sf = int(snowflake_fetch_data(monitoring_queries["DTS_DC_DIMENTIONAL_TABLE"]["DIM_SHAZAM_CHART_ARTIST_ITUNES_SF"])[0]['COUNT']) dim_shahzam_charts_artist_itunes_pg = int(postgres_fetch_data(monitoring_queries["DTS_DC_DIMENTIONAL_TABLE"]["DIM_SHAZAM_CHART_ARTIST_ITUNES_PG"])[0]['count']) print(f"Check 5 DTS_DC_DIMENTIONAL_TABLE DIM_SHAZAM_CHARTS_SF - {dim_shahzam_charts_sf}") print(f"Check 5 DTS_DC_DIMENTIONAL_TABLE DIM_SHAZAM_CHARTS_PG - {dim_shahzam_charts_pg}") print(f"Check 5 DTS_DC_DIMENTIONAL_TABLE DIM_SHAZAM_CITY_SF - {dim_shahzam_city_sf}") print(f"Check 5 DTS_DC_DIMENTIONAL_TABLE DIM_SHAZAM_CITY_PG - {dim_shahzam_city_pg}") print(f"Check 5 DTS_DC_DIMENTIONAL_TABLE DIM_SHAZAM_TRACK_SF - {dim_shahzam_track_sf}") print(f"Check 5 DTS_DC_DIMENTIONAL_TABLE DIM_SHAZAM_TRACK_PG - {dim_shahzam_track_pg}") print(f"Check 5 DTS_DC_DIMENTIONAL_TABLE SHAZAM_CHART_TRACK_SUMMARY_REPORT_SF - {dim_shahzam_track_summary_report_sf}") print(f"Check 5 DTS_DC_DIMENTIONAL_TABLE SHAZAM_CHART_TRACK_SUMMARY_REPORT_PG - {dim_shahzam_track_summary_report_pg}") print(f"Check 5 DTS_DC_DIMENTIONAL_TABLE DIM_SHAZAM_CHART_ARTIST_SF - {dim_shahzam_charts_artist_sf}") print(f"Check 5 DTS_DC_DIMENTIONAL_TABLE DIM_SHAZAM_CHART_ARTIST_PG - {dim_shahzam_charts_artist_pg}") print(f"Check 5 DTS_DC_DIMENTIONAL_TABLE DIM_SHAZAM_CHART_ARTIST_ITUNES_SF - {dim_shahzam_charts_artist_itunes_sf}") print(f"Check 5 DTS_DC_DIMENTIONAL_TABLE DIM_SHAZAM_CHART_ARTIST_ITUNES_PG - {dim_shahzam_charts_artist_itunes_pg}") full_log["DTS_DC_DIMENTIONAL_TABLE_DIM_SHAZAM_CHARTS"] = {'DIM_SHAZAM_CHARTS_SF': dim_shahzam_charts_sf, 'DIM_SHAZAM_CHARTS_PG': dim_shahzam_charts_pg} full_log["DTS_DC_DIMENTIONAL_TABLE_DIM_SHAZAM_CITY"] = {'DIM_SHAZAM_CITY_SF': dim_shahzam_city_sf, 'DIM_SHAZAM_CITY_PG': dim_shahzam_city_pg} full_log["DTS_DC_DIMENTIONAL_TABLE_DIM_SHAZAM_TRACK"] = {'DIM_SHAZAM_TRACK_SF': dim_shahzam_track_sf, 'DIM_SHAZAM_TRACK_PG': dim_shahzam_track_pg} full_log["DTS_DC_DIMENTIONAL_TABLE_SHAZAM_CHART_TRACK_SUMMARY_REPORT"] = {'SHAZAM_CHART_TRACK_SUMMARY_REPORT_SF': dim_shahzam_track_summary_report_sf, 'SHAZAM_CHART_TRACK_SUMMARY_REPORT_PG': dim_shahzam_track_summary_report_pg} full_log["DTS_DC_DIMENTIONAL_TABLE_DIM_SHAZAM_CHART_ARTIST"] = {'DIM_SHAZAM_CHART_ARTIST_SF': dim_shahzam_charts_artist_sf, 'DIM_SHAZAM_CHART_ARTIST_PG': dim_shahzam_charts_artist_pg} full_log["DTS_DC_DIMENTIONAL_TABLE_DIM_SHAZAM_CHART_ARTIST_ITUNES"] = {'DIM_SHAZAM_CHART_ARTIST_ITUNES_SF': dim_shahzam_charts_artist_itunes_sf, 'DIM_SHAZAM_CHART_ARTIST_ITUNES_PG': dim_shahzam_charts_artist_itunes_pg} if dim_shahzam_charts_sf != dim_shahzam_charts_pg: result_dict["DTS_DC_DIMENTIONAL_TABLE_DIM_SHAZAM_CHARTS"] = full_log["DTS_DC_DIMENTIONAL_TABLE_DIM_SHAZAM_CHARTS"] if dim_shahzam_city_sf != dim_shahzam_city_pg: result_dict["DTS_DC_DIMENTIONAL_TABLE_DIM_SHAZAM_CITY"] = full_log["DTS_DC_DIMENTIONAL_TABLE_DIM_SHAZAM_CITY"] if dim_shahzam_track_sf != dim_shahzam_track_pg: result_dict["DTS_DC_DIMENTIONAL_TABLE_DIM_SHAZAM_TRACK"] = full_log["DTS_DC_DIMENTIONAL_TABLE_DIM_SHAZAM_TRACK"] if dim_shahzam_track_summary_report_sf != dim_shahzam_track_summary_report_pg: result_dict["DTS_DC_DIMENTIONAL_TABLE_SHAZAM_CHART_TRACK_SUMMARY_REPORT"] = full_log["DTS_DC_DIMENTIONAL_TABLE_SHAZAM_CHART_TRACK_SUMMARY_REPORT"] if dim_shahzam_charts_artist_sf != dim_shahzam_charts_artist_pg: result_dict["DTS_DC_DIMENTIONAL_TABLE_DIM_SHAZAM_CHART_ARTIST"] = full_log["DTS_DC_DIMENTIONAL_TABLE_DIM_SHAZAM_CHART_ARTIST"] if dim_shahzam_charts_artist_itunes_sf != dim_shahzam_charts_artist_itunes_pg: result_dict["DTS_DC_DIMENTIONAL_TABLE_DIM_SHAZAM_CHART_ARTIST_ITUNES"] = full_log["DTS_DC_DIMENTIONAL_TABLE_DIM_SHAZAM_CHART_ARTIST_ITUNES"] # Check 6 DTS_DC_FACT_TABLE fact_shahzam_charts_sf = int(snowflake_fetch_data(monitoring_queries["DTS_DC_FACT_TABLE"]["FACT_SHAZAM_CHART_SF"])[0]['COUNT']) fact_shahzam_charts_pg = int(postgres_fetch_data(monitoring_queries["DTS_DC_FACT_TABLE"]["FACT_SHAZAM_CHART_PG"])[0]['count']) fact_shahzam_track_total_sf = int(snowflake_fetch_data(monitoring_queries["DTS_DC_FACT_TABLE"]["FACT_SHAZAM_TRACK_TOTAL_SHAZAMS_SF"])[0]['COUNT']) fact_shahzam_track_total_pg = int(postgres_fetch_data(monitoring_queries["DTS_DC_FACT_TABLE"]["FACT_SHAZAM_TRACK_TOTAL_SHAZAMS_PG"])[0]['count']) print(f"Check 6 DTS_DC_FACT_TABLE FACT_SHAZAM_CHART_SF - {fact_shahzam_charts_sf}") print(f"Check 6 DTS_DC_FACT_TABLE FACT_SHAZAM_CHART_PG - {fact_shahzam_charts_pg}") print(f"Check 6 DTS_DC_FACT_TABLE FACT_SHAZAM_TRACK_TOTAL_SHAZAMS_SF - {fact_shahzam_track_total_sf}") print(f"Check 6 DTS_DC_FACT_TABLE FACT_SHAZAM_TRACK_TOTAL_SHAZAMS_PG - {fact_shahzam_track_total_pg}") full_log["DTS_DC_FACT_TABLE_FACT_SHAZAM_CHART"] = {'FACT_SHAZAM_CHART_SF': fact_shahzam_charts_sf, 'FACT_SHAZAM_CHART_PG': fact_shahzam_charts_pg} full_log["DTS_DC_FACT_TABLE_FACT_SHAZAM_TRACK_TOTAL_SHAZAMS"] = {'FACT_SHAZAM_TRACK_TOTAL_SHAZAMS_SF': fact_shahzam_track_total_sf, 'FACT_SHAZAM_TRACK_TOTAL_SHAZAMS_PG': fact_shahzam_track_total_pg} if fact_shahzam_charts_sf != fact_shahzam_charts_pg: result_dict["DTS_DC_FACT_TABLE_FACT_SHAZAM_CHART"] = full_log["DTS_DC_FACT_TABLE_FACT_SHAZAM_CHART"] if fact_shahzam_track_total_sf != fact_shahzam_track_total_pg: result_dict["DTS_DC_FACT_TABLE_FACT_SHAZAM_TRACK_TOTAL_SHAZAMS"] = full_log["DTS_DC_FACT_TABLE_FACT_SHAZAM_TRACK_TOTAL_SHAZAMS"] return result_dict, full_log ################# run the all steps ###################### def run_steps() -> dict: message_dict = {'name':"[PROD] Shazam 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=0)).strftime("%Y-%m-%d") result_dict_123, monitoring_123_log = monitoring_123() # result_dict_4, monitoring_4_log = monitoring_4() result_dict_56, monitoring_56_log = monitoring_56() # message_dict['full_log'] = {**monitoring_123_log, **monitoring_4_log, **monitoring_56_log} message_dict['full_log'] = {**monitoring_123_log, **monitoring_56_log} # if result_dict_4 or result_dict_56 or result_dict_123: if result_dict_56 or result_dict_123: postgres_update_failure_count(failure_count_increase) if result_dict_123: # Data is partially synced with Chartmetric message_dict['status'] = "FAILED" message_dict['info'] = "Data is partially synced with Chartmetric" message_dict['details'] = result_dict_123 elif result_dict_56: # Data is synced with Chartmetric and partially available in API message_dict['status'] = "FAILED" message_dict['info'] = "Data is synced with Chartmetric and partially available in API" message_dict['details'] = result_dict_56 # if result_dict_4: # message_dict['details'].update(result_dict_4) else: message_dict['status'] = "OK" message_dict['info'] = "Data is synced with Chartmetric and available in API" postgres_update_failure_count(failure_count_to_zero) 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/Shazam_Charts_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()