#! /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"] # 'as count' column name was added in all count() entities in all SQL queries to ease further processing. Please also add it in case of query updates. monitoring_queries = { "DS_CHARTMETRICK_VS_RAW": { "DS_CHARTMETRICK": "SELECT count(*) as count, 'source' as tbl FROM DS_CHARTMETRIC.RAW_DATA.YOUTUBE WHERE id is not null union all SELECT count(*) as count, 'destination' as tbl FROM DELPHI_CHARTMETRIC_RAW.CHARTS_YOUTUBE.YOUTUBE;", "RAW": "SELECT count(distinct POSITION,COUNTRY,TIMESTP) as count, 'source' as tbl FROM DS_CHARTMETRIC.RAW_DATA.YOUTUBE_VIDEO_TRENDING_CHART WHERE TIMESTP <= dateadd(dd, -1, current_date()) union all SELECT count(distinct POSITION,COUNTRY,TIMESTP) as count, 'destination' as tbl FROM DELPHI_CHARTMETRIC_RAW.CHARTS_YOUTUBE.YOUTUBE_VIDEO_TRENDING_CHART WHERE TIMESTP <= dateadd(dd, -1, current_date());" }, "RAW_VS_MAIN": { "RAW": "SELECT count(*) as count, 'source' as tbl FROM DELPHI_CHARTMETRIC_RAW.CHARTS_YOUTUBE.YOUTUBE union all SELECT count(*) as count, 'destination' as tbl FROM DELPHI_CHARTMETRIC_MAIN.CHARTS_YOUTUBE.DIM_YOUTUBE_CHART_META;", "MAIN": "SELECT count(distinct POSITION,COUNTRY,TIMESTP) as count, 'source' as tbl FROM DELPHI_CHARTMETRIC_RAW.CHARTS_YOUTUBE.YOUTUBE_VIDEO_TRENDING_CHART union all SELECT count(*) as count, 'destination' as tbl FROM DELPHI_CHARTMETRIC_MAIN.CHARTS_YOUTUBE.FACT_YOUTUBE_CHART_DAILY;" }, "MAIN_VS_ETL": { "MAIN": "SELECT count(*) as count, 'source' as tbl FROM DELPHI_CHARTMETRIC_MAIN.CHARTS_YOUTUBE.FACT_YOUTUBE_CHART_DAILY union all SELECT count(*) as count, 'destination' as tbl FROM DELPHI_CHARTMETRIC_ETL.CHARTS_YOUTUBE.FACT_YOUTUBE_CHART_DAILY;", "ETL": "SELECT count(*) as count, 'source' as tbl FROM DELPHI_CHARTMETRIC_MAIN.CHARTS_YOUTUBE.DIM_YOUTUBE_CHART_META union all SELECT count(*) as count, 'destination' as tbl FROM DELPHI_CHARTMETRIC_ETL.CHARTS_YOUTUBE.DIM_YOUTUBE_CHART_META;" }, "DATA_TRANSFER_STATUS_UOW": "SELECT uow_meta.unit_of_work.uow_id,uow_meta.unit_of_work.status,uow_meta.audit_log.extra_data,uow_meta.audit_log.new_status FROM uow_meta.unit_of_work INNER JOIN uow_meta.audit_log ON uow_meta.unit_of_work.uow_id = uow_meta.audit_log.uow_id WHERE status = 'IMPORT_FAILED' AND project = 'cm_charts_yt' AND new_status = 'IMPORT_FAILED' AND uow_meta.audit_log.created_at >= current_date - 2", "DATA_TRANSFER_STATUS_DATA_CONSISTENCY_DIMENTIONAL_TABLE": { "DIM_YOUTUBE_CHART_META": "SELECT count(*) as count FROM DELPHI_CHARTMETRIC_ETL.CHARTS_YOUTUBE.DIM_YOUTUBE_CHART_META", "DIM_YOUTUBE_VIDEO": "SELECT count(*) as count FROM chartmetric.dim_youtube_video" }, "DATA_TRANSFER_STATUS_DATA_CONSISTENCY_FACT_TABLE": { "FACT_YOUTUBE_CHART_DAILY": "SELECT count(*) as count FROM DELPHI_CHARTMETRIC_ETL.CHARTS_YOUTUBE.FACT_YOUTUBE_CHART_DAILY where VIDEO_ID in (select VIDEO_ID from DELPHI_CHARTMETRIC_ETL.CHARTS_YOUTUBE.DIM_YOUTUBE_CHART_META)", "FACT_YOUTUBE_VIDEO_TRENDING_CHART": "SELECT count(*) as count FROM chartmetric.fact_youtube_video_trending_chart" } } report_steps = { 'step_1': "SELECT sum(*) FROM (SELECT count(*) FROM DS_CHARTMETRIC.RAW_DATA.YOUTUBE WHERE id is not null union all SELECT count(*) FROM DS_CHARTMETRIC.RAW_DATA.YOUTUBE_VIDEO_TRENDING_CHART)", 'step_2': "SELECT sum(*) FROM (SELECT count(*) FROM DELPHI_CHARTMETRIC_ETL.CHARTS_YOUTUBE.DIM_YOUTUBE_CHART_META union all SELECT count(*) FROM DELPHI_CHARTMETRIC_ETL.CHARTS_YOUTUBE.FACT_YOUTUBE_CHART_DAILY)", 'step_3': "SELECT sum(*) FROM (SELECT count(*) FROM chartmetric.dim_youtube_video union all SELECT count(*) FROM chartmetric.fact_youtube_video_trending_chart)" } failure_count_increase = "UPDATE failed_count SET failed_count = failed_count + 1, last = 'FAILED', last_updated = NOW() WHERE count_id = 'yttch';" failure_count_to_zero = "UPDATE failed_count SET failed_count = 0, last = 'OK', last_updated = NOW() WHERE count_id = 'yttch';" failure_count = "select failed_count from failed_count where count_id='yttch';" 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] YouTube Trending 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 postgres_fetch_data_local(sql_query) -> list: result_list = [] 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) 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 monitoring_123() -> dict: result_dict = {} full_log = {} # Check 1 DS_CHARTMETRICK_VS_RAW tmp_fetched_ds_chartmetrick = snowflake_fetch_data(monitoring_queries["DS_CHARTMETRICK_VS_RAW"]["DS_CHARTMETRICK"]) ds_chartmetrick_src = int(tmp_fetched_ds_chartmetrick[0]['COUNT']) ds_chartmetrick_dst = int(tmp_fetched_ds_chartmetrick[1]['COUNT']) tmp_fetched_raw = snowflake_fetch_data(monitoring_queries["DS_CHARTMETRICK_VS_RAW"]["RAW"]) raw_src = int(tmp_fetched_raw[0]['COUNT']) raw_dst = int(tmp_fetched_raw[1]['COUNT']) print(f"Check 1 DS_CHARTMETRICK_VS_RAW YOUTUBE - {tmp_fetched_ds_chartmetrick}") print(f"Check 1 DS_CHARTMETRICK_VS_RAW YOUTUBE_VIDEO_TRENDING_CHART - {tmp_fetched_raw}") full_log["DS_CHARTMETRICK_VS_RAW_1"] = {'DS_CHARTMETRICK source': ds_chartmetrick_src, 'DS_CHARTMETRICK destination': ds_chartmetrick_dst} full_log["DS_CHARTMETRICK_VS_RAW_2"] = {'RAW source': raw_src, 'RAW destination': raw_dst} if ds_chartmetrick_src != ds_chartmetrick_dst: result_dict["DS_CHARTMETRICK_VS_RAW_1"] = full_log["DS_CHARTMETRICK_VS_RAW_1"] if raw_src != raw_dst: result_dict["DS_CHARTMETRICK_VS_RAW_2"] = full_log["DS_CHARTMETRICK_VS_RAW_2"] # Check 2 RAW_VS_MAIN tmp_fetched_raw = snowflake_fetch_data(monitoring_queries["RAW_VS_MAIN"]["RAW"]) raw_src = int(tmp_fetched_raw[0]['COUNT']) raw_dst = int(tmp_fetched_raw[1]['COUNT']) tmp_fetched_main = snowflake_fetch_data(monitoring_queries["RAW_VS_MAIN"]["MAIN"]) main_src = int(tmp_fetched_main[0]['COUNT']) main_dst = int(tmp_fetched_main[1]['COUNT']) print(f"Check 2 RAW_VS_MAIN YOUTUBE - {tmp_fetched_raw}") print(f"Check 2 RAW_VS_MAIN YOUTUBE_VIDEO_TRENDING_CHART - {tmp_fetched_main}") full_log["RAW_VS_MAIN_1"] = {'RAW source': raw_src, 'RAW destination': raw_dst} full_log["RAW_VS_MAIN_2"] = {'MAIN source': main_src, 'MAIN destination': main_dst} if raw_src != raw_dst: result_dict["RAW_VS_MAIN_1"] = full_log["RAW_VS_MAIN_1"] if main_src != main_dst: result_dict["RAW_VS_MAIN_2"] = full_log["RAW_VS_MAIN_2"] # Check 3 MAIN_VS_ETL tmp_fetched_main = snowflake_fetch_data(monitoring_queries["MAIN_VS_ETL"]["MAIN"]) main_src = int(tmp_fetched_main[0]['COUNT']) main_dst = int(tmp_fetched_main[1]['COUNT']) tmp_fetched_etl = snowflake_fetch_data(monitoring_queries["MAIN_VS_ETL"]["ETL"]) etl_src = int(tmp_fetched_etl[0]['COUNT']) etl_dst = int(tmp_fetched_etl[1]['COUNT']) print(f"Check 3 MAIN_VS_ETL FACT_YOUTUBE_CHART_DAILY - {tmp_fetched_main}") print(f"Check 3 MAIN_VS_ETL DIM_YOUTUBE_CHART_META - {tmp_fetched_etl}") full_log["MAIN_VS_ETL_1"] = {'MAIN source': main_src, 'MAIN destination': main_dst} full_log["MAIN_VS_ETL_2"] = {'ETL source': etl_src, 'ETL destination': etl_dst} if main_src != main_dst: result_dict["MAIN_VS_ETL_1"] = full_log["MAIN_VS_ETL_1"] if etl_src != etl_dst: result_dict["MAIN_VS_ETL_2"] = full_log["MAIN_VS_ETL_2"] return result_dict, full_log def monitoring_4() -> dict: result_dict = {} full_log = {} # Check 4 DATA_TRANSFER_STATUS_UOW data_transfer_status_uow = postgres_fetch_data(monitoring_queries["DATA_TRANSFER_STATUS_UOW"]) print(f"Check 4 DATA_TRANSFER_STATUS_UOW - {data_transfer_status_uow}") full_log["DATA_TRANSFER_STATUS_UOW"] = data_transfer_status_uow if data_transfer_status_uow: result_dict["DATA_TRANSFER_STATUS_UOW"] = data_transfer_status_uow return result_dict, full_log def monitoring_56() -> dict: result_dict = {} full_log = {} # Check 5 DATA_TRANSFER_STATUS_DATA_CONSISTENCY_DIMENTIONAL_TABLE dim_youtube_chart_meta = int(snowflake_fetch_data(monitoring_queries["DATA_TRANSFER_STATUS_DATA_CONSISTENCY_DIMENTIONAL_TABLE"]["DIM_YOUTUBE_CHART_META"])[0]['COUNT']) dim_youtube_video = int(postgres_fetch_data(monitoring_queries["DATA_TRANSFER_STATUS_DATA_CONSISTENCY_DIMENTIONAL_TABLE"]["DIM_YOUTUBE_VIDEO"])[0]['count']) print(f"Check 5 DATA_TRANSFER_STATUS_DATA_CONSISTENCY_DIMENTIONAL_TABLE SNOWFLAKE - {dim_youtube_chart_meta}") print(f"Check 5 DATA_TRANSFER_STATUS_DATA_CONSISTENCY_DIMENTIONAL_TABLE POSTGRES- {dim_youtube_video}") full_log["DATA_TRANSFER_STATUS_DATA_CONSISTENCY_DIMENTIONAL_TABLE"] = {'DIM_YOUTUBE_CHART_META': dim_youtube_chart_meta, 'DIM_YOUTUBE_VIDEO': dim_youtube_video} if dim_youtube_chart_meta != dim_youtube_video: result_dict["DATA_TRANSFER_STATUS_DATA_CONSISTENCY_DIMENTIONAL_TABLE"] = full_log["DATA_TRANSFER_STATUS_DATA_CONSISTENCY_DIMENTIONAL_TABLE"] # Check 6 DATA_TRANSFER_STATUS_DATA_CONSISTENCY_FACT_TABLE fact_youtube_chart_daily = int(snowflake_fetch_data(monitoring_queries["DATA_TRANSFER_STATUS_DATA_CONSISTENCY_FACT_TABLE"]["FACT_YOUTUBE_CHART_DAILY"])[0]['COUNT']) fact_youtube_video_trending_chart = int(postgres_fetch_data(monitoring_queries["DATA_TRANSFER_STATUS_DATA_CONSISTENCY_FACT_TABLE"]["FACT_YOUTUBE_VIDEO_TRENDING_CHART"])[0]['count']) print(f"Check 6 DATA_TRANSFER_STATUS_DATA_CONSISTENCY_FACT_TABLE SNOWFLAKE - {fact_youtube_chart_daily}") print(f"Check 6 DATA_TRANSFER_STATUS_DATA_CONSISTENCY_FACT_TABLE POSTGRES- {fact_youtube_video_trending_chart}") full_log["DATA_TRANSFER_STATUS_DATA_CONSISTENCY_FACT_TABLE"] = {'FACT_YOUTUBE_CHART_DAILY': fact_youtube_chart_daily, 'FACT_YOUTUBE_VIDEO_TRENDING_CHART': fact_youtube_video_trending_chart} if fact_youtube_chart_daily != fact_youtube_video_trending_chart: result_dict["DATA_TRANSFER_STATUS_DATA_CONSISTENCY_FACT_TABLE"] = full_log["DATA_TRANSFER_STATUS_DATA_CONSISTENCY_FACT_TABLE"] return result_dict, full_log ################# run the all steps ###################### def run_steps() -> dict: message_dict = {'name':"[PROD] YouTube Trending Charts",'status':"",'date':"", 'info':"", 'details':"", 'full_log':"", 'failure_count':""} # 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") count_yttch = int(postgres_fetch_data_local(failure_count)[0]['failed_count']) 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} 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 message_dict['step'] = "1" message_dict['failure_count'] = count_yttch 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 message_dict['step'] = "2" message_dict['failure_count'] = count_yttch # if result_dict_4 and message_dict['details']: # 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/Youtube_Trending_charts_hourly_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()