#! /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] DAPD Apple/Spotify Playlists 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() -> dict: result_dict = {} full_log = {} ### Check 1 ################################################## for query_name, query in sql_queries.pg_queries_1.items(): tmp_check = postgres_fetch_data_slz(f"""{query}""") if tmp_check: full_log[f"{query_name}"] = {} result_dict[f"{query_name}"] = {} print(f"{query_name}: - {tmp_check}") for item in tmp_check: full_log[f"{query_name}"].update({item['report_date'].strftime("%Y-%m-%d"): item['units_list']}) print(f"{item['report_date']} - {item['units_list']}") result_dict[f"{query_name}"].update({item['report_date'].strftime("%Y-%m-%d"): item['units_list']}) else: print(f"{query_name} returned 0 results") full_log[f"{query_name}"] = {f"OK: Query returned 0 rows"} return result_dict, full_log def check_2_4() -> dict: result_dict = {} full_log = {} ### Check 2 ################################################## for query_name, query in sql_queries.pg_queries_2.items(): tmp_check = postgres_fetch_data_slz(f"""{query}""") if tmp_check: full_log[f"{query_name}"] = {} result_dict[f"{query_name}"] = {} print(f"{query_name}: - {tmp_check}") for item in tmp_check: full_log[f"{query_name}"].update({item['report_date'].strftime("%Y-%m-%d"): item['units_list']}) print(f"{item['report_date']} - {item['units_list']}") result_dict[f"{query_name}"].update({item['report_date'].strftime("%Y-%m-%d"): item['units_list']}) else: print(f"{query_name} returned 0 results") full_log[f"{query_name}"] = {f"OK: Query returned 0 rows"} ### Check 3 ################################################## for query_name, query in sql_queries.pg_queries_3.items(): tmp_check = postgres_fetch_data_workflowdb(f"""{query}""") if tmp_check: full_log[f"{query_name}"] = {} result_dict[f"{query_name}"] = {} print(f"{query_name}: - {tmp_check}") for item in tmp_check: full_log[f"{query_name}"].update({item['added_at'].strftime("%Y-%m-%d"): f"{item['data_source']} {item['entity']}" }) print(f"{item['added_at']} - {item['data_source']}, {item['entity']}") result_dict[f"{query_name}"].update({item['added_at'].strftime("%Y-%m-%d"): f"{item['data_source']} {item['entity']}" }) else: print(f"{query_name} returned 0 results") full_log[f"{query_name}"] = {f"OK: Query returned 0 rows"} ### Check 4 ################################################## for query_name, query in sql_queries.pg_queries_4.items(): tmp_check = postgres_fetch_data_etldb(f"""{query}""") if tmp_check: full_log[f"{query_name}"] = {} result_dict[f"{query_name}"] = {} print(f"{query_name}: - {tmp_check}") for item in tmp_check: full_log[f"{query_name}"].update({item['added_at'].strftime("%Y-%m-%d"): f"dsp_id:{item['dsp_id']}, entity:{item['entity']}" }) print(f"{item['added_at']} - {item['dsp_id']}, {item['entity']}") result_dict[f"{query_name}"].update({item['added_at'].strftime("%Y-%m-%d"): f"dsp_id:{item['dsp_id']}, entity:{item['entity']}" }) else: print(f"{query_name} returned 0 results") full_log[f"{query_name}"] = {f"OK: Query returned 0 rows"} return result_dict, full_log def check_5_7() -> dict: result_dict = {} full_log = {} ### Check 5 ################################################## for query_name, query in sql_queries.pg_queries_5.items(): tmp_check = postgres_fetch_data_main(f"""{query}""") if tmp_check: full_log[f"{query_name}"] = {} result_dict[f"{query_name}"] = {} print(f"{query_name}: - {tmp_check}") for item in tmp_check: full_log[f"{query_name}"].update({item['updated_at'].strftime("%Y-%m-%d"): f"status:{item['status']}, project:{item['project']}, uow_id:{item['unit_of_work_id']}" }) print(f"{item['updated_at']} - status:{item['status']}, project:{item['project']}, uow_id:{item['unit_of_work_id']}") result_dict[f"{query_name}"].update({item['updated_at'].strftime("%Y-%m-%d"): f"status:{item['status']}, project:{item['project']}, uow_id:{item['unit_of_work_id']}" }) else: print(f"{query_name} returned 0 results") full_log[f"{query_name}"] = {f"OK: Query returned 0 rows"} ### Check 7 ################################################## for query_name, query in sql_queries.pg_queries_7.items(): tmp_check = postgres_fetch_data_main(f"""{query}""") ### tmp_check = [{'table': 'test_table_1', 'count': 0},{'table': 'test_table_2', 'count': 0}] if tmp_check: full_log[f"{query_name}"] = {} result_dict_tmp = {} result_dict_tmp[f"{query_name}"] = {} print(f"{query_name}: - {tmp_check}") for item in tmp_check: full_log[f"{query_name}"].update({item['table']: item['count']}) if item['count'] != 0: result_dict[f"{query_name}"] = {} print(f"{item['table']} - {item['count']}") result_dict_tmp[f"{query_name}"].update({item['table']: item['count']}) result_dict[f"{query_name}"] = result_dict_tmp[f"{query_name}"] else: print(f"{query_name} OK: Count in every row == 0") full_log[f"{query_name}"] = {f"OK: Count in every row == 0"} return result_dict, full_log def check_9_57() -> dict: result_dict = {} full_log = {} ### Check 9-23 ################################################## for query_name, query in sql_queries.sf_queries_9_23.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 or 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 24-41 ################################################## for query_name, query in sql_queries.sf_queries_24_41.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 42-57 ################################################## ### UPDATE QUERY 55 WHEN FIXED for query_name, query in sql_queries.sf_queries_42_57.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 61-67 ################################################## for check_name, check_data in sql_queries.queries_61_67.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 abs(sf_check_count - pg_check_count) > 1000: result_dict[f"{check_name}"] = full_log[f"{check_name}"] return result_dict, full_log def check_58_60() -> dict: result_dict = {} full_log = {} ### Check 58-60 ################################################## for query_name, query in sql_queries.sf_queries_58_60.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 or 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}"] return result_dict, full_log ################# run the all steps ###################### def run_steps() -> dict: message_dict = {'name':"[PROD] DAPD Apple/Spotify Playlists Monitoring",'status':"",'date':"", 'info':"", 'details':"", 'full_log':""} # set date according to sql searching date, currently it's `today` message_dict["date"] = datetime.datetime.now().strftime("%Y-%m-%d") result_dict_1, monitoring_1_log = check_1() result_dict_2_4, monitoring_2_4_log = check_2_4() result_dict_5_7, monitoring_5_7_log = check_5_7() result_dict_9_57, monitoring_9_57_log = check_9_57() result_dict_58_60, monitoring_58_60_log = check_58_60() message_dict['full_log'] = {**monitoring_1_log, **monitoring_2_4_log, **monitoring_5_7_log, **monitoring_9_57_log, **monitoring_58_60_log} message_dict['details'] = {**result_dict_1, **result_dict_2_4, **result_dict_5_7, **result_dict_9_57, **result_dict_58_60} if result_dict_1: message_dict['status'] = "FAILED" message_dict['info'] = "Data is partially synced with Apollo" elif result_dict_2_4: message_dict['status'] = "FAILED" message_dict['info'] = "Data is synced with Apollo, but partially ingested from APIs" elif result_dict_5_7: message_dict['status'] = "FAILED" message_dict['info'] = "Data is ingested from APIs, but partially available in Delphi API" elif result_dict_9_57: message_dict['status'] = "FAILED" message_dict['info'] = "Data is available in Delphi API, but partially available in Delphi Snowflake" elif result_dict_58_60: message_dict['status'] = "FAILED" message_dict['info'] = "Data is available in Delphi Snowflake, but notifications are partially generated" else: message_dict['status'] = "OK" message_dict['info'] = "Data is available in Delphi API and Snowflake; Notifications are generated" 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/dapd_spotify_apple_hourly_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()