#! /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 = 'lnkfr';" failure_count_to_zero = "UPDATE failed_count SET failed_count = 0, last = 'OK', last_updated = NOW() WHERE count_id = 'lnkfr';" 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] Linkfire" 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 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 check_1() -> dict: result_dict = {} full_log = {} ### Check 1 SLZ status ################################################ 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_6() -> dict: result_dict = {} full_log = {} ### Check 2 SLZ → EXP - Data transfer status - Import ################################################ tmp_sf_slz_exp_check_2 = snowflake_fetch_data(sql_queries.sf_slz_exp_check_2) print(f"SLZ_TO_EXP_DTS_IMPORT_CHECK_2: - {tmp_sf_slz_exp_check_2}") #Check if tmp_sf_slz_exp_check_2: result_dict["SLZ_TO_EXP_DTS_IMPORT_CHECK_2"] = {} full_log["SLZ_TO_EXP_DTS_IMPORT_CHECK_2"] = {} for item in tmp_sf_slz_exp_check_2: full_log["SLZ_TO_EXP_DTS_IMPORT_CHECK_2"].update({item['REPORT_DATE'].strftime("%Y-%m-%d"): item['MISSING_REPORTS'].lower()}) print(f"{item['REPORT_DATE']} - {item['MISSING_REPORTS']}") result_dict["SLZ_TO_EXP_DTS_IMPORT_CHECK_2"].update({item['REPORT_DATE'].strftime("%Y-%m-%d"): item['MISSING_REPORTS'].lower()}) else: full_log["SLZ_TO_EXP_DTS_IMPORT_CHECK_2"] = {'SLZ_TO_EXP_DTS_IMPORT_CHECK_2': 'Query returned 0 results'} ### Check 3 EXP → LINKFIRE - Data transfer status - import status successful ################################################ tmp_sf_exp_lf_success_check_3 = snowflake_fetch_data(sql_queries.sf_exp_lf_success_check_3) exp_lf_success_check_3_sf = int(tmp_sf_exp_lf_success_check_3[0]['count']) print(f"EXP_TO_LF_DTS_IMPORT_SUCCESS_CHECK_3: - {exp_lf_success_check_3_sf} success operations") full_log["EXP_TO_LF_DTS_IMPORT_SUCCESS_CHECK_3"] = {} #Check if exp_lf_success_check_3_sf == 0: result_dict["EXP_TO_LF_DTS_IMPORT_SUCCESS_CHECK_3"] = {} result_dict["EXP_TO_LF_DTS_IMPORT_SUCCESS_CHECK_3"] = {'EXP_TO_LF_DTS_IMPORT_SUCCESS_CHECK_3': 'Query returned 0 count for completed operations'} full_log["EXP_TO_LF_DTS_IMPORT_SUCCESS_CHECK_3"] = result_dict["EXP_TO_LF_DTS_IMPORT_SUCCESS_CHECK_3"] else: full_log["EXP_TO_LF_DTS_IMPORT_SUCCESS_CHECK_3"] = {'EXP_TO_LF_DTS_IMPORT_SUCCESS_CHECK_3': 'Query returned completed operations'} ### Check 4 EXP → LINKFIRE - Data transfer status - import status failed ################################################ tmp_sf_exp_lf_failed_check_4 = snowflake_fetch_data(sql_queries.sf_exp_lf_failed_check_4) exp_lf_failed_check_4_sf = int(tmp_sf_exp_lf_failed_check_4[0]['count']) print(f"EXP_TO_LF_DTS_IMPORT_FAILED_CHECK_4: - {exp_lf_failed_check_4_sf} failed operations") full_log["EXP_TO_LF_DTS_IMPORT_FAILED_CHECK_4"] = {} #Check if exp_lf_failed_check_4_sf > 0: result_dict["EXP_TO_LF_DTS_IMPORT_FAILED_CHECK_4"] = {} result_dict["EXP_TO_LF_DTS_IMPORT_FAILED_CHECK_4"] = {'EXP_TO_LF_DTS_IMPORT_FAILED_CHECK_4': 'Query returned incompleted operations'} full_log["EXP_TO_LF_DTS_IMPORT_FAILED_CHECK_4"] = result_dict["EXP_TO_LF_DTS_IMPORT_FAILED_CHECK_4"] else: full_log["EXP_TO_LF_DTS_IMPORT_FAILED_CHECK_4"] = {'EXP_TO_LF_DTS_IMPORT_FAILED_CHECK_4': 'Query returned 0 count for incomplete operations'} ### Check 5 EXP → LINKFIRE - Data transfer status - Data consistency - FACT_EVENT_FUNNEL ################################################ tmp_dts_dc_check_5 = snowflake_fetch_data(sql_queries.sf_dts_dc_check_5) dts_dc_check_5_src = int(tmp_dts_dc_check_5[0]['count']) dts_dc_check_5_dest = int(tmp_dts_dc_check_5[1]['count']) print(f"CHECK_5_DTS_DC_FACT_EVENT_FUNNEL_SRC: - {dts_dc_check_5_src}") print(f"CHECK_5_DTS_DC_FACT_EVENT_FUNNEL_DEST: - {dts_dc_check_5_dest}") full_log["CHECK_5_DTS_DC_FACT_EVENT_FUNNEL"] = {'DTS_DC_FACT_EVENT_FUNNEL_SRC': dts_dc_check_5_src, 'DTS_DC_FACT_EVENT_FUNNEL_DEST': dts_dc_check_5_dest} if dts_dc_check_5_src != dts_dc_check_5_dest: result_dict["CHECK_5_DTS_DC_FACT_EVENT_FUNNEL"] = full_log["CHECK_5_DTS_DC_FACT_EVENT_FUNNEL"] ### Check 6 EXP → LINKFIRE - Data transfer status - Data consistency - DIM_TERRITORY ################################################ tmp_dts_dc_check_6 = snowflake_fetch_data(sql_queries.sf_dts_dc_check_6) dts_dc_check_6_src = int(tmp_dts_dc_check_6[0]['count']) dts_dc_check_6_dest = int(tmp_dts_dc_check_6[1]['count']) print(f"CHECK_6_DTS_DC_DIM_TERRITORY_SRC: - {dts_dc_check_6_src}") print(f"CHECK_6_DTS_DC_DIM_TERRITORY_DEST: - {dts_dc_check_6_dest}") full_log["CHECK_6_DTS_DC_DIM_TERRITORY"] = {'DTS_DC_DIM_TERRITORY_SRC': dts_dc_check_6_src, 'DTS_DC_DIM_TERRITORY_DEST': dts_dc_check_6_dest} if dts_dc_check_6_src != dts_dc_check_6_dest: result_dict["CHECK_6_DTS_DC_DIM_TERRITORY"] = {'DTS_DC_DIM_TERRITORY_SRC': dts_dc_check_6_src, 'DTS_DC_DIM_TERRITORY_DEST': dts_dc_check_6_dest} ### Check 7 EXP → LINKFIRE - Data transfer status - Data consistency - DIM_LINKFIRE_LINK ################################################ tmp_dts_dc_check_7 = snowflake_fetch_data(sql_queries.sf_dts_dc_check_7) dts_dc_check_7_src = int(tmp_dts_dc_check_7[0]['count']) dts_dc_check_7_dest = int(tmp_dts_dc_check_7[1]['count']) print(f"CHECK_7_DTS_DC_DIM_LINKFIRE_LINK_SRC: - {dts_dc_check_7_src}") print(f"CHECK_7_DTS_DC_DIM_LINKFIRE_LINK_DEST: - {dts_dc_check_7_dest}") full_log["CHECK_7_DTS_DC_DIM_LINKFIRE_LINK"] = {'DTS_DC_DIM_LINKFIRE_LINK_SRC': dts_dc_check_7_src, 'DTS_DC_DIM_LINKFIRE_LINK_DEST': dts_dc_check_7_dest} if dts_dc_check_7_src != dts_dc_check_7_dest: result_dict["CHECK_7_DTS_DC_DIM_LINKFIRE_LINK"] = full_log["CHECK_7_DTS_DC_DIM_LINKFIRE_LINK"] return result_dict, full_log def check_10_to_16() -> dict: result_dict = {} full_log = {} ### Check 10 LINKFIRE → S3 - Data transfer status - export status successful ################################################ tmp_lf_s3_success_check_10 = snowflake_fetch_data(sql_queries.sf_lf_s3_success_check_10) lf_s3_success_check_10_sf = int(tmp_lf_s3_success_check_10[0]['count']) print(f"LF_TO_S3_DTS_IMPORT_SUCCESS_CHECK_10: - {lf_s3_success_check_10_sf} success operations") full_log["LF_TO_S3_DTS_IMPORT_SUCCESS_CHECK_10"] = {} #Check if lf_s3_success_check_10_sf == 0: result_dict["LF_TO_S3_DTS_IMPORT_SUCCESS_CHECK_10"] = {} result_dict["LF_TO_S3_DTS_IMPORT_SUCCESS_CHECK_10"] = {'LF_TO_S3_DTS_IMPORT_SUCCESS_CHECK_10': 'Query returned 0 count for completed operations'} full_log["LF_TO_S3_DTS_IMPORT_SUCCESS_CHECK_10"] = result_dict["LF_TO_S3_DTS_IMPORT_SUCCESS_CHECK_10"] else: full_log["LF_TO_S3_DTS_IMPORT_SUCCESS_CHECK_10"] = {'LF_TO_S3_DTS_IMPORT_SUCCESS_CHECK_10': 'Query returned completed operations'} ### Check 11 LINKFIRE → S3 - Data transfer status - export status failed ################################################ tmp_sf_lf_s3_failed_check_11 = snowflake_fetch_data(sql_queries.sf_lf_s3_failed_check_11) lf_s3_failed_check_11_sf = int(tmp_sf_lf_s3_failed_check_11[0]['count']) print(f"LF_TO_S3_DTS_IMPORT_FAILED_CHECK_11: - {lf_s3_failed_check_11_sf} failed operations") full_log["LF_TO_S3_DTS_IMPORT_FAILED_CHECK_11"] = {} #Check if lf_s3_failed_check_11_sf > 0: result_dict["LF_TO_S3_DTS_IMPORT_FAILED_CHECK_11"] = {} result_dict["LF_TO_S3_DTS_IMPORT_FAILED_CHECK_11"] = {'LF_TO_S3_DTS_IMPORT_FAILED_CHECK_11': 'Query returned incompleted operations'} full_log["LF_TO_S3_DTS_IMPORT_FAILED_CHECK_11"] = result_dict["LF_TO_S3_DTS_IMPORT_FAILED_CHECK_11"] else: full_log["LF_TO_S3_DTS_IMPORT_FAILED_CHECK_11"] = {'LF_TO_S3_DTS_IMPORT_FAILED_CHECK_11': 'Query returned 0 count for incomplete operations'} ### Check 13 Data transfer status - Data consistency - DIM_LINKFIRE_LINK ################################################ tmp_sf_dts_dc_check_13 = snowflake_fetch_data(sql_queries.sf_dts_dc_check_13) tmp_pg_dts_dc_check_13 = postgres_fetch_data_main(sql_queries.pg_dts_dc_check_13) dts_dc_check_13_sf = int(tmp_sf_dts_dc_check_13[0]['count']) dts_dc_check_13_pg = int(tmp_pg_dts_dc_check_13[0]['count']) print(f"CHECK_13_DTS_DC_DIM_LINKFIRE_LINK_SF: - {dts_dc_check_13_sf}") print(f"CHECK_13_DTS_DC_DIM_LINKFIRE_LINK_PG: - {dts_dc_check_13_pg}") full_log["CHECK_13_DTS_DC_DIM_LINKFIRE_LINK"] = {'DTS_DC_DIM_LINKFIRE_LINK_SF': dts_dc_check_13_sf, 'DTS_DC_DIM_LINKFIRE_LINK_PG': dts_dc_check_13_pg} if dts_dc_check_13_sf != dts_dc_check_13_pg: result_dict["CHECK_13_DTS_DC_DIM_LINKFIRE_LINK"] = full_log["CHECK_13_DTS_DC_DIM_LINKFIRE_LINK"] ### Check 14 Data transfer status - Data consistency - FACT_EVENT_FUNNEL ################################################ tmp_sf_dts_dc_check_14 = snowflake_fetch_data(sql_queries.sf_dts_dc_check_14) tmp_pg_dts_dc_check_14 = postgres_fetch_data_main(sql_queries.pg_dts_dc_check_14) dts_dc_check_14_sf = int(tmp_sf_dts_dc_check_14[0]['count']) dts_dc_check_14_pg = int(tmp_pg_dts_dc_check_14[0]['count']) print(f"CHECK_14_DTS_DC_FACT_EVENT_FUNNEL_SF: - {dts_dc_check_14_sf}") print(f"CHECK_14_DTS_DC_FACT_EVENT_FUNNEL_PG: - {dts_dc_check_14_pg}") full_log["CHECK_14_DTS_DC_FACT_EVENT_FUNNEL"] = {'DTS_DC_FACT_EVENT_FUNNEL_SF': dts_dc_check_14_sf, 'DTS_DC_FACT_EVENT_FUNNEL_PG': dts_dc_check_14_pg} if dts_dc_check_14_sf != dts_dc_check_14_pg: result_dict["CHECK_14_DTS_DC_FACT_EVENT_FUNNEL"] = full_log["CHECK_14_DTS_DC_FACT_EVENT_FUNNEL"] ### Check 15 Data transfer status - Data consistency - DIM_REFERRER ################################################ tmp_sf_dts_dc_check_15 = snowflake_fetch_data(sql_queries.sf_dts_dc_check_15) tmp_pg_dts_dc_check_15 = postgres_fetch_data_main(sql_queries.pg_dts_dc_check_15) dts_dc_check_15_sf = int(tmp_sf_dts_dc_check_15[0]['count']) dts_dc_check_15_pg = int(tmp_pg_dts_dc_check_15[0]['count']) print(f"CHECK_15_DTS_DC_DIM_REFERRER_SF: - {dts_dc_check_15_sf}") print(f"CHECK_15_DTS_DC_DIM_REFERRER_PG: - {dts_dc_check_15_pg}") full_log["CHECK_15_DTS_DC_DIM_REFERRER"] = {'DTS_DC_DIM_REFERRER_SF': dts_dc_check_15_sf, 'DTS_DC_DIM_REFERRER_PG': dts_dc_check_15_pg} if dts_dc_check_15_sf != dts_dc_check_15_pg: result_dict["CHECK_15_DTS_DC_DIM_REFERRER"] = full_log["CHECK_15_DTS_DC_DIM_REFERRER"] ### Check 16 Data transfer status - Data consistency - DIM_CAMPAIGN_LINK ################################################ tmp_sf_dts_dc_check_16 = snowflake_fetch_data(sql_queries.sf_dts_dc_check_16) tmp_pg_dts_dc_check_16 = postgres_fetch_data_main(sql_queries.pg_dts_dc_check_16) dts_dc_check_16_sf = int(tmp_sf_dts_dc_check_16[0]['count']) dts_dc_check_16_pg = int(tmp_pg_dts_dc_check_16[0]['count']) print(f"CHECK_16_DTS_DC_DIM_CAMPAIGN_LINK_SF: - {dts_dc_check_16_sf}") print(f"CHECK_16_DTS_DC_DIM_CAMPAIGN_LINK_PG: - {dts_dc_check_16_pg}") full_log["CHECK_16_DTS_DC_DIM_CAMPAIGN_LINK"] = {'DTS_DC_DIM_CAMPAIGN_LINK_SF': dts_dc_check_16_sf, 'DTS_DC_DIM_CAMPAIGN_LINK_PG': dts_dc_check_16_pg} if dts_dc_check_16_sf != dts_dc_check_16_pg: result_dict["CHECK_16_DTS_DC_DIM_CAMPAIGN_LINK"] = full_log["CHECK_16_DTS_DC_DIM_CAMPAIGN_LINK"] return result_dict, full_log ################# run the all steps ###################### def run_steps() -> dict: message_dict = {'name':"[PROD] Linkfire",'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=2)).strftime("%Y-%m-%d") result_dict_1, monitoring_1_log = check_1() result_dict_2_to_6, monitoring_2_to_6_log = check_2_to_6() result_dict_10_to_16, monitoring_10_to_16_log = check_10_to_16() message_dict['full_log'] = {**monitoring_1_log, **monitoring_2_to_6_log, **monitoring_10_to_16_log} message_dict['details'] = {**result_dict_1, **result_dict_2_to_6, **result_dict_10_to_16} if result_dict_2_to_6 or result_dict_1 or result_dict_10_to_16: 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_6: # 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_10_to_16: 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/Linkfire_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()