import psycopg2 from contextlib import closing import os import csv import json import re import logging from datetime import datetime, timedelta # # POSTGRES credentials from environment variables # database = os.environ.get("POSTGRES_DATABASE") # user = os.environ.get("POSTGRES_USER") # password = os.environ.get("POSTGRES_PASSWORD") # host = os.environ.get("POSTGRES_HOST") # port = os.environ.get("POSTGRES_PORT") # POSTGRES credentials from json file with open('config.json', 'r') as f: config = json.load(f) sme_database = config['db_sme']["database"] sme_user = config['db_sme']["user"] sme_password = config['db_sme']["password"] sme_host = config['db_sme']["host"] sme_port = config['db_sme']["port"] report_log_database = config['db_daily_report']["database"] report_log_user = config['db_daily_report']["user"] report_log_password = config['db_daily_report']["password"] report_log_host = config['db_daily_report']["host"] report_log_port = config['db_daily_report']["port"] # execute sql query and create eventual csv file def fetch_csv(date: str, query: str, to_file: str) -> None: if date == datetime.now().strftime("%Y-%m-%d"): copy_sql = f"COPY ({query}) TO STDOUT DELIMITER ',' CSV HEADER;" else: query_adjusted_by_date = re.sub("and uow.report_date >= current_date - 14", f"and uow.report_date = '{date}'", query) copy_sql = f"COPY ({query_adjusted_by_date}) TO STDOUT DELIMITER ',' CSV HEADER;" with closing(psycopg2.connect(database=sme_database, user=sme_user, password=sme_password, host=sme_host, port=sme_port)) as connection: with connection.cursor() as cursor: with open(f"{to_file}", 'w') as csv_file: cursor.copy_expert(copy_sql, csv_file) def select(query): result = [] try: with closing(psycopg2.connect(database=report_log_database, user=report_log_user, password=report_log_password, host=report_log_host, port=report_log_port)) as connection: with connection.cursor() as cursor: cursor.execute(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.append(tmp) except Exception as e: logging.error(e) return result def update(source, issue_date, status): query = f"""update audit_log set status = '{status}', completed_at = now() where source = '{source}' and issue_date = '{issue_date}';""" try: with closing(psycopg2.connect(database=report_log_database, user=report_log_user, password=report_log_password, host=report_log_host, port=report_log_port)) as connection: with connection.cursor() as cursor: cursor.execute(query) connection.commit() except Exception as e: logging.error(e) def insert(source, issue_date, status): query = f"""insert into audit_log(source,created_at,issue_date,status) values ('{source}',now(),'{issue_date}','{status}');""" try: with closing(psycopg2.connect(database=report_log_database, user=report_log_user, password=report_log_password, host=report_log_host, port=report_log_port)) as connection: with connection.cursor() as cursor: cursor.execute(query) connection.commit() except Exception as e: logging.error(e) def etl_processing(file: str, date: str, sources_list: dict) -> dict: csv_file = open(str(file), "r") dict_reader = csv.DictReader(csv_file) csv_list_dict = [dict(x) for x in dict_reader] csv_file.close() exception_conditions = { 'youtubereporting': { 'etl_status': ['MIN_COMPLETE', 'COMPLETE'], 'report_names': ['youtubereporting_streams_isrc_date_day', 'youtubereporting_demographics_isrc_date_day', 'youtubereporting_streams_video_id_date_day', 'youtubereporting_demographics_video_id_date_day'] } } status_messages = { 'OK': "Ingested and ETL Processed*", 'FAILED': "Some UOWs aren't COMPLETED*", 'ABSENT': "DSP is absent in query result*" } # get eventual DSP list without duplicates from csv sources_list_from_csv = set(x['sources'] for x in csv_list_dict) # get uow list that have wrong 'etl_status' uow_incompleted = [x for x in csv_list_dict if x['etl_status'] != 'COMPLETE'] # get uow list with all statuses for today for additional custom check uow_all = [x for x in csv_list_dict] # prepare empty result message structure result_data = {x: {'date': '', 'status': '', 'message': '', 'details': []} for x in sources_list.keys()} # fill result message with data according to checking rules for source in result_data.keys(): # use offset for data_sources only if the passed `date` to function is `today` if date == datetime.now().strftime("%Y-%m-%d"): source_days_offset = int(sources_list[source]['days_offset']) else: source_days_offset = 0 result_data[source]['date'] = (datetime.strptime(date, '%Y-%m-%d') - timedelta(days=source_days_offset)).strftime("%Y-%m-%d") result_data[source]['details'] = [x for x in uow_incompleted if x['sources'] == source and x['report_date'] == result_data[source]['date']] # here we check if DSP presents in query result and then continue with other checks and fill result message if source in sources_list_from_csv: if len(result_data[source]['details']) == 0: result_data[source]['status'] = "OK" # additional check for looking dsp->report_names with acceptable etl_status according to exception_conditions elif source in exception_conditions.keys(): result_data[source]['details'] = [x for x in uow_all if x['sources'] == source and x['report_date'] == result_data[source]['date']] # generate details with all uow for additional checks checking_list = [] for element in result_data[source]['details']: if element['report_name'] in exception_conditions[source]['report_names'] and element['etl_status'] in exception_conditions[source]['etl_status']: checking_list.append(element) if len(checking_list) == len(exception_conditions[source]['report_names']): result_data[source]['status'] = "OK" else: result_data[source]['status'] = "FAILED" else: result_data[source]['status'] = "FAILED" else: result_data[source]['status'] = "ABSENT" # set up the message according to the source checking status result_data[source]['message'] = status_messages[result_data[source]['status']] return result_data if __name__ == '__main__': pass