#! /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 from datetime import date, timedelta 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_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"] 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_workflowdb = config['postgres_workflowdb']["database"] db_user_workflowdb = config['postgres_workflowdb']["user"] db_password_workflowdb = config['postgres_workflowdb']["password"] db_host_workflowdb = config['postgres_workflowdb']["host"] db_port_workflowdb = config['postgres_workflowdb']["port"] db_name_etldb = config['postgres_etldb']["database"] db_user_etldb = config['postgres_etldb']["user"] db_password_etldb = config['postgres_etldb']["password"] db_host_etldb = config['postgres_etldb']["host"] db_port_etldb = config['postgres_etldb']["port"] 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] Amazon Charts Monitoring" 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_etldb(sql_query) -> list: result_list = [] with closing(psycopg2.connect(database=db_name_etldb, user=db_user_etldb, password=db_password_etldb, host=db_host_etldb, port=db_port_etldb)) 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_workflowdb(sql_query) -> list: result_list = [] with closing(psycopg2.connect(database=db_name_workflowdb, user=db_user_workflowdb, password=db_password_workflowdb, host=db_host_workflowdb, port=db_port_workflowdb)) 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_5() -> dict: result_dict = {} full_log = {} ### Check 1-5 ################################################## for query_name, query in sql_queries.sf_queries_1_5.items(): tmp_check = snowflake_fetch_data(f"""{query}""") full_log[f"{query_name}"] = {} if tmp_check: src_count = tmp_check[0]['count'] dest_count = tmp_check[1]['count'] print(f"{query_name} - src: {src_count}, dest: {dest_count}") full_log[f"{query_name}"].update({f"src": src_count}) full_log[f"{query_name}"].update({f"dest": dest_count}) if src_count != dest_count: result_dict[f"{query_name}"] = full_log[f"{query_name}"] else: full_log[f"{query_name}"] = {f"FAILED: Query returned nothing"} result_dict[f"{query_name}"] = full_log[f"{query_name}"] return result_dict, full_log def check_6_27() -> dict: result_dict = {} full_log = {} ### Check 6-11,14-16,18-20 ################################################## for query_name, query in sql_queries.sf_queries_6_20.items(): tmp_check = snowflake_fetch_data(f"""{query}""") full_log[f"{query_name}"] = {} if tmp_check: src_count = tmp_check[0]['count'] dest_count = tmp_check[1]['count'] print(f"{query_name} - src: {src_count}, dest: {dest_count}") full_log[f"{query_name}"].update({f"src": src_count}) full_log[f"{query_name}"].update({f"dest": dest_count}) if src_count != dest_count: result_dict[f"{query_name}"] = full_log[f"{query_name}"] else: full_log[f"{query_name}"] = {f"FAILED: Query returned nothing"} result_dict[f"{query_name}"] = full_log[f"{query_name}"] ### Check 12,13,17 ################################################## for query_name, query in sql_queries.sf_queries_12_13_17.items(): tmp_check = snowflake_fetch_data(f"""{query}""") full_log[f"{query_name}"] = {} if tmp_check: count = tmp_check[0]['count'] print(f"{query_name} count - {count}") full_log[f"{query_name}"].update({f"count": count}) if count != 0: result_dict[f"{query_name}"] = full_log[f"{query_name}"] else: full_log[f"{query_name}"] = {f"FAILED: Query returned nothing"} result_dict[f"{query_name}"] = full_log[f"{query_name}"] ### Check 21-26 ################################################## for check_name, check_data in sql_queries.sf_queries_21_26.items(): sf_query = check_data.get('sf_check_' + check_name[-2:]) pg_query = check_data.get('pg_check_' + check_name[-2:]) print(f"sf_query_{check_name}: {sf_query}") print(f"pg_query_{check_name}: {pg_query}") tmp_sf_check = snowflake_fetch_data(sf_query) tmp_pg_check = postgres_fetch_data_main(pg_query) if tmp_sf_check and tmp_pg_check: sf_check_count = int(tmp_sf_check[0]['count']) pg_check_count = int(tmp_pg_check[0]['count']) print(f"sf_{check_name}: - {sf_check_count}") print(f"pg_{check_name}: - {pg_check_count}") full_log[f"{check_name}"] = {f"sf_{check_name}": sf_check_count, f"pg_{check_name}": pg_check_count} if sf_check_count != pg_check_count: result_dict[f"{check_name}"] = full_log[f"{check_name}"] ### Check 27 ################################################## for query_name, query in sql_queries.pg_queries_27.items(): tmp_check = postgres_fetch_data_main(f"""{query}""") if tmp_check: print(f"{query_name}: - {tmp_check}") full_log[f"{query_name}"] = {} result_dict[f"{query_name}"] = {} for item in tmp_check: uow = item['uow_id'] status = item['status'] extra = item['extra_data'] new_status = item['new_status'] full_log[f"{query_name}"][f"{uow}"] = {} full_log[f"{query_name}"][f"{uow}"].update({"uow_id":uow, "status": status, "extra_data": extra, "new_status": new_status}) result_dict[f"{query_name}"][f"{uow}"] = {} result_dict[f"{query_name}"][f"{uow}"].update({"uow_id":uow, "status": status, "extra_data": extra, "new_status": new_status}) else: full_log[f"{query_name}"] = {f"OK: All UoWs are COMPLETED"} return result_dict, full_log ################# run the all steps ###################### def run_steps() -> dict: message_dict = {'name':"[PROD] Amazon Charts Monitoring",'status':"",'date':"", 'info':"", 'details':"", 'full_log':""} # set date according to sql searching date, currently it's `today` message_dict["date"] = message_dict["date"] = (datetime.datetime.now() - datetime.timedelta(days=1)).strftime("%Y-%m-%d") result_dict_1_5, monitoring_1_5_log = check_1_5() result_dict_6_27, monitoring_6_27_log = check_6_27() message_dict['full_log'] = {**monitoring_1_5_log, **monitoring_6_27_log} message_dict['details'] = {**result_dict_1_5, **result_dict_6_27} sorted_full_log = dict(sorted(message_dict['full_log'].items())) message_dict['full_log'] = sorted_full_log sorted_details = dict(sorted(message_dict['details'].items())) message_dict['details'] = sorted_details if result_dict_1_5: message_dict['status'] = "FAILED" message_dict['info'] = "Data is partially synced with Chartmetric" elif result_dict_6_27: message_dict['status'] = "FAILED" message_dict['info'] = "Data is synced with Chartmetric and partially available in API" else: message_dict['status'] = "OK" message_dict['info'] = "Data is synced with Chartmetric and 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/amazon_charts_daily_log.json', 'w') as json_file: json.dump(result, json_file, indent=4, default=str) if arg_list.recipients: recipients = arg_list.recipients[0] send_mail(recipients=recipients, message=json.dumps(result, indent=4, default=str), status=result['status']) if arg_list.print: print(json.dumps(result, indent=4, default=str)) else: parser.print_help()