#! /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"] 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 = 'tktkdsb';" failure_count_to_zero = "UPDATE failed_count SET failed_count = 0, last = 'OK', last_updated = NOW() WHERE count_id = 'tktkdsb';" 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] TikTok Dashboard" 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 check_1_to_14() -> dict: result_dict = {} full_log = {} ### Check 1 EXPLORATION → RAW TIKTOK_TRENDS_DAILY (expected result: Counts match inside every step) ################################################ tmp_sf_check_1 = snowflake_fetch_data(sql_queries.sf_check_1) sf_check_1_src = int(tmp_sf_check_1[0]['count']) sf_check_1_dest = int(tmp_sf_check_1[1]['count']) print(f"SRC_CHECK_1: - {sf_check_1_src}") print(f"DEST_CHECK_1: - {sf_check_1_dest}") full_log["CHECK_1"] = {'SRC_CHECK_1': sf_check_1_src, 'DEST_CHECK_1': sf_check_1_dest} if sf_check_1_src != sf_check_1_dest: result_dict["CHECK_1"] = full_log["CHECK_1"] ### Check 2 RAW -> MAIN DIM_TIKTOK_TOP_CHART (expected result: Counts match inside every step) ################################################ tmp_sf_check_2 = snowflake_fetch_data(sql_queries.sf_check_2) sf_check_2_src = int(tmp_sf_check_2[0]['count']) sf_check_2_dest = int(tmp_sf_check_2[1]['count']) print(f"SRC_CHECK_2: - {sf_check_2_src}") print(f"DEST_CHECK_2: - {sf_check_2_dest}") full_log["CHECK_2"] = {'SRC_CHECK_2': sf_check_2_src, 'DEST_CHECK_2': sf_check_2_dest} if sf_check_2_src != sf_check_2_dest: result_dict["CHECK_2"] = full_log["CHECK_2"] ### Check 3 RAW -> MAIN DIM_TIKTOK_TOP_SONG (expected result: Counts match inside every step) ################################################ tmp_sf_check_3 = snowflake_fetch_data(sql_queries.sf_check_3) sf_check_3_src = int(tmp_sf_check_3[0]['count']) sf_check_3_dest = int(tmp_sf_check_3[1]['count']) print(f"SRC_CHECK_3: - {sf_check_3_src}") print(f"DEST_CHECK_3: - {sf_check_3_dest}") full_log["CHECK_3"] = {'SRC_CHECK_3': sf_check_3_src, 'DEST_CHECK_3': sf_check_3_dest} if sf_check_3_src != sf_check_3_dest: result_dict["CHECK_3"] = full_log["CHECK_3"] ### Check 4 MAIN -> MAIN FACT_TIKTOK_TOP_CHART_ISRC (expected result: Counts match inside every step) ################################################ tmp_sf_check_4 = snowflake_fetch_data(sql_queries.sf_check_4) sf_check_4_src = int(tmp_sf_check_4[0]['count']) sf_check_4_dest = int(tmp_sf_check_4[1]['count']) print(f"SRC_CHECK_4: - {sf_check_4_src}") print(f"DEST_CHECK_4: - {sf_check_4_dest}") full_log["CHECK_4"] = {'SRC_CHECK_4': sf_check_4_src, 'DEST_CHECK_4': sf_check_4_dest} if sf_check_4_src != sf_check_4_dest: result_dict["CHECK_4"] = full_log["CHECK_4"] ### Check 5 MAIN -> MAIN FACT_TIKTOK_TOP_CHART_ISRC (expected result: Counts match inside every step) ################################################ tmp_sf_check_5 = snowflake_fetch_data(sql_queries.sf_check_5) sf_check_5_src = int(tmp_sf_check_5[0]['count']) sf_check_5_dest = int(tmp_sf_check_5[1]['count']) print(f"SRC_CHECK_5: - {sf_check_5_src}") print(f"DEST_CHECK_5: - {sf_check_5_dest}") full_log["CHECK_5"] = {'SRC_CHECK_5': sf_check_5_src, 'DEST_CHECK_5': sf_check_5_dest} if sf_check_5_src != sf_check_5_dest: result_dict["CHECK_5"] = full_log["CHECK_5"] ### Check 6 MAIN -> ETL DIM_TIKTOK_TOP_CHART (expected result: Counts match inside every step) ################################################ tmp_sf_check_6 = snowflake_fetch_data(sql_queries.sf_check_6) sf_check_6_src = int(tmp_sf_check_6[0]['count']) sf_check_6_dest = int(tmp_sf_check_6[1]['count']) print(f"SRC_CHECK_6: - {sf_check_6_src}") print(f"DEST_CHECK_6: - {sf_check_6_dest}") full_log["CHECK_6"] = {'SRC_CHECK_6': sf_check_6_src, 'DEST_CHECK_6': sf_check_6_dest} if sf_check_6_src != sf_check_6_dest: result_dict["CHECK_6"] = full_log["CHECK_6"] ### Check 7 MAIN -> ETL FACT_TIKTOK_TOP_CHART_ISRC (expected result: Counts match inside every step) ################################################ tmp_sf_check_7 = snowflake_fetch_data(sql_queries.sf_check_7) sf_check_7_src = int(tmp_sf_check_7[0]['count']) sf_check_7_dest = int(tmp_sf_check_7[1]['count']) print(f"SRC_CHECK_7: - {sf_check_7_src}") print(f"DEST_CHECK_7: - {sf_check_7_dest}") full_log["CHECK_7"] = {'SRC_CHECK_7': sf_check_7_src, 'DEST_CHECK_7': sf_check_7_dest} if sf_check_7_src != sf_check_7_dest: result_dict["CHECK_7"] = full_log["CHECK_7"] ### Check 8 MAIN -> ETL FACT_TIKTOK_TOP_CHART_ISRC_LIFETIME (expected result: Counts match inside every step) ################################################ tmp_sf_check_8 = snowflake_fetch_data(sql_queries.sf_check_8) sf_check_8_src = int(tmp_sf_check_8[0]['count']) sf_check_8_dest = int(tmp_sf_check_8[1]['count']) print(f"SRC_CHECK_8: - {sf_check_8_src}") print(f"DEST_CHECK_8: - {sf_check_8_dest}") full_log["CHECK_8"] = {'SRC_CHECK_8': sf_check_8_src, 'DEST_CHECK_8': sf_check_8_dest} if sf_check_8_src != sf_check_8_dest: result_dict["CHECK_8"] = full_log["CHECK_8"] ### Check 9 MAIN -> ETL DIM_TIKTOK_TOP_SONG (expected result: Counts match inside every step) ################################################ tmp_sf_check_9 = snowflake_fetch_data(sql_queries.sf_check_9) sf_check_9_src = int(tmp_sf_check_9[0]['count']) sf_check_9_dest = int(tmp_sf_check_9[1]['count']) print(f"SRC_CHECK_9: - {sf_check_9_src}") print(f"DEST_CHECK_9: - {sf_check_9_dest}") full_log["CHECK_9"] = {'SRC_CHECK_9': sf_check_9_src, 'DEST_CHECK_9': sf_check_9_dest} if sf_check_9_src != sf_check_9_dest: result_dict["CHECK_9"] = full_log["CHECK_9"] ### Check 11 ETL -> PG DIM_TIKTOK_TOP_CHART (expected result: Counts match between steps) ################################################ tmp_sf_check_11 = snowflake_fetch_data(sql_queries.sf_check_11) tmp_pg_check_11 = postgres_fetch_data_main(sql_queries.pg_check_11) sf_check_11_count = int(tmp_sf_check_11[0]['count']) pg_check_11_count = int(tmp_pg_check_11[0]['count']) print(f"SRC_CHECK_11: - {sf_check_11_count}") print(f"DEST_CHECK_11: - {pg_check_11_count}") full_log["CHECK_11"] = {'SRC_CHECK_11': sf_check_11_count, 'DEST_CHECK_11': pg_check_11_count} if sf_check_11_count != pg_check_11_count: result_dict["CHECK_11"] = full_log["CHECK_11"] ### Check 12 ETL -> PG DIM_TIKTOK_TOP_SONG (expected result: Counts match between steps) ################################################ tmp_sf_check_12 = snowflake_fetch_data(sql_queries.sf_check_12) tmp_pg_check_12 = postgres_fetch_data_main(sql_queries.pg_check_12) sf_check_12_count = int(tmp_sf_check_12[0]['count']) pg_check_12_count = int(tmp_pg_check_12[0]['count']) print(f"SRC_CHECK_12: - {sf_check_12_count}") print(f"DEST_CHECK_12: - {pg_check_12_count}") full_log["CHECK_12"] = {'SRC_CHECK_12': sf_check_12_count, 'DEST_CHECK_12': pg_check_12_count} if sf_check_12_count != pg_check_12_count: result_dict["CHECK_12"] = full_log["CHECK_12"] ### Check 13 ETL -> PG fact_tiktok_top_chart_isrc (expected result: Counts match between steps) ################################################ tmp_sf_check_13 = snowflake_fetch_data(sql_queries.sf_check_13) tmp_pg_check_13 = postgres_fetch_data_main(sql_queries.pg_check_13) sf_check_13_count = int(tmp_sf_check_13[0]['count']) pg_check_13_count = int(tmp_pg_check_13[0]['count']) print(f"SRC_CHECK_13: - {sf_check_13_count}") print(f"DEST_CHECK_13: - {pg_check_13_count}") full_log["CHECK_13"] = {'SRC_CHECK_13': sf_check_13_count, 'DEST_CHECK_13': pg_check_13_count} if sf_check_13_count != pg_check_13_count: result_dict["CHECK_13"] = full_log["CHECK_13"] ### Check 14 ETL -> PG fact_tiktok_top_chart_isrc_lifetime (expected result: Counts match between steps) ################################################ tmp_sf_check_14 = snowflake_fetch_data(sql_queries.sf_check_14) tmp_pg_check_14 = postgres_fetch_data_main(sql_queries.pg_check_14) sf_check_14_count = int(tmp_sf_check_14[0]['count']) pg_check_14_count = int(tmp_pg_check_14[0]['count']) print(f"SRC_CHECK_14: - {sf_check_14_count}") print(f"DEST_CHECK_14: - {pg_check_14_count}") full_log["CHECK_14"] = {'SRC_CHECK_14': sf_check_14_count, 'DEST_CHECK_14': pg_check_14_count} if sf_check_14_count != pg_check_14_count: result_dict["CHECK_14"] = full_log["CHECK_14"] return result_dict, full_log ################# run the all steps ###################### def run_steps() -> dict: message_dict = {'name':"[PROD] TikTok Dashboard",'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_to_14, monitoring_1_to_14_log = check_1_to_14() message_dict['full_log'] = {**monitoring_1_to_14_log} message_dict['details'] = {**result_dict_1_to_14} if result_dict_1_to_14: postgres_update_failure_count(failure_count_increase) message_dict['status'] = "FAILED" message_dict['info'] = "Counts do not match inside the steps below" else: postgres_update_failure_count(failure_count_to_zero) message_dict['status'] = "OK" message_dict['info'] = "Counts match inside every step" 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/TikTok_Dashboard_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()