import datetime import json import os from dataclasses import dataclass import slack from sme_logger import get_logger from aws_utils import AwsUtils from const import ENV logger = get_logger('CHARTMETRIC_BACKFILL', os.environ.get("ENVIRONMENT", ENV).lower() != ENV) @dataclass class ChartmetricSnowflakeUtils: chartmetric_constants = json.loads(AwsUtils().get_secret( required_secret=f"delphi/{ENV}/chartmetric_constants")) slack_channel = chartmetric_constants.get("CHARTMETRIC_SLACK_CHANNEL") slack_token = chartmetric_constants.get("CHARTMETRIC_SLACK_TOKEN") slack_client = slack.WebClient(token=slack_token) exception_records = [] # Snowflake to Bigtable migration part begins - - - - - - > def update_process_table(self, *, PROCESS_ID=None, RUN_DATE=None, STATUS=None, START_TIME=None, END_TIME=None, RUN_TYPE=None, SOURCE=None, SINCE=None, UNTIL=None, is_updating=False): if SINCE is None: SINCE = str(datetime.datetime.now() - datetime.timedelta(days=7)).split(" ")[0] if is_updating: update_process_table_query = f""" UPDATE "{self.chartmetric_constants.get("SNOWFLAKE_DATABASE")}"."{self.chartmetric_constants.get("SNOWFLAKE_SCHEMA")}".CHARTMETRIC_PROCESS SET STATUS = 'COMPLETE', END_TIME = '{END_TIME}' WHERE PROCESS_ID = '{PROCESS_ID}' AND RUN_DATE = '{RUN_DATE}' AND START_TIME = '{START_TIME}' AND STATUS = 'IN_PROGRESS'; """ else: update_process_table_query = f""" INSERT INTO "{self.chartmetric_constants.get("SNOWFLAKE_DATABASE")}" ."{self.chartmetric_constants.get("SNOWFLAKE_SCHEMA")}" .CHARTMETRIC_PROCESS ( PROCESS_ID, RUN_DATE, STATUS, START_TIME, RUN_TYPE,--SCHEDULED/MANUAL SOURCE, SINCE, UNTIL ) VALUES ('{PROCESS_ID}','{RUN_DATE}', '{STATUS}', '{START_TIME}', '{RUN_TYPE}', '{SOURCE}','{SINCE}', '{UNTIL}'); """ return update_process_table_query def generate_data_checker_query(self, *, source): data_checker_query = f""" CREATE OR REPLACE TABLE "{self.chartmetric_constants.get("SNOWFLAKE_DATABASE")}"."{self.chartmetric_constants.get("SNOWFLAKE_SCHEMA")}".SME_{source}_STAT AS SELECT * FROM "{self.chartmetric_constants.get("SNOWFLAKE_DATABASE")}"."{self.chartmetric_constants.get("SNOWFLAKE_SCHEMA")}".V_SME_{source}_STAT; """ return data_checker_query def generate_data_fetch_query(self, *, source: str, since: str = None, until: str = None): if since is not None and until is not None: data_fetch_query = f""" SELECT * from "{self.chartmetric_constants.get("SNOWFLAKE_DATABASE")}"."{self.chartmetric_constants.get("SNOWFLAKE_SCHEMA")}".SME_{source}_STAT stat WHERE report_date BETWEEN NVL('{since}', dateadd(week, -1, '{until}')) and '{until}' """ else: data_fetch_query = f""" select stat.* from "{self.chartmetric_constants.get("SNOWFLAKE_DATABASE")}"."{self.chartmetric_constants.get("SNOWFLAKE_SCHEMA")}".SME_{source}_STAT stat LEFT JOIN (SELECT a.chartmetric_id, max(report_date) as max_2, max({source.lower()}_max) as max_1 FROM "{self.chartmetric_constants.get("SNOWFLAKE_DATABASE")}"."{self.chartmetric_constants.get("SNOWFLAKE_SCHEMA")}".SME_{source}_STAT a LEFT JOIN "{self.chartmetric_constants.get("SNOWFLAKE_DATABASE")}"."{self.chartmetric_constants.get("SNOWFLAKE_SCHEMA")}".CHARTMETRIC_STATUS b ON b.chartmetric_id=a.chartmetric_id group by a.chartmetric_id) a ON stat.chartmetric_id=a.chartmetric_id WHERE stat.report_date between NVL(a.max_1, dateadd(week, -1, current_date())) and a.max_2; """ return data_fetch_query def generate_status_table_query(self, *, source, since: str = None, until: str = None, schedule_value): update_status_table_query = None if since is None or until is None: update_status_table_query = f""" MERGE INTO "{self.chartmetric_constants.get("SNOWFLAKE_DATABASE")}"."{self.chartmetric_constants.get("SNOWFLAKE_SCHEMA")}".CHARTMETRIC_STATUS AS target USING ( SELECT a.chartmetric_id ,max(report_date) max_date ,listagg(report_date, ',') within GROUP ( ORDER BY report_date ) AS processed_date FROM "{self.chartmetric_constants.get("SNOWFLAKE_DATABASE")}"."{self.chartmetric_constants.get("SNOWFLAKE_SCHEMA")}"."SME_{source}_STAT" AS a ,( SELECT a.chartmetric_id ,max(report_date) AS max_2 ,max({source}_max) AS max_1 FROM "{self.chartmetric_constants.get("SNOWFLAKE_DATABASE")}"."{self.chartmetric_constants.get("SNOWFLAKE_SCHEMA")}"."SME_{source}_STAT" a LEFT JOIN "{self.chartmetric_constants.get("SNOWFLAKE_DATABASE")}" ."{self.chartmetric_constants.get("SNOWFLAKE_SCHEMA")}" .CHARTMETRIC_STATUS b ON b.chartmetric_id = a.chartmetric_id GROUP BY a.chartmetric_id ) AS c WHERE a.chartmetric_id = c.chartmetric_id AND report_date BETWEEN NVL(max_1,dateadd(week, -1, current_date())) AND max_2 GROUP BY a.chartmetric_id ) AS b ON target.chartmetric_id = b.chartmetric_id AND target.run_date = CURRENT_DATE () AND NVL(target.schedule, 0) > {schedule_value} WHEN MATCHED THEN UPDATE SET target.{source}_max = b.max_date ,target.{source} = b.processed_date WHEN NOT MATCHED THEN INSERT ( chartmetric_id ,run_date ,STATUS ,{source}_max ,{source} ,schedule ) VALUES ( b.chartmetric_id ,CURRENT_DATE () ,'BACKFILL' ,b.max_date ,b.processed_date ,{int(schedule_value) + 1} ); """ if since is not None and until is not None: update_status_table_query = f""" MERGE INTO "{self.chartmetric_constants.get("SNOWFLAKE_DATABASE")}"."{self.chartmetric_constants.get("SNOWFLAKE_SCHEMA")}".chartmetric_status AS target using ( SELECT a.chartmetric_id , Max(report_date) max_date , Listagg(report_date, ',') within GROUP ( ORDER BY report_date ) AS processed_date FROM "{self.chartmetric_constants.get("SNOWFLAKE_DATABASE")}"."{self.chartmetric_constants.get("SNOWFLAKE_SCHEMA")}"."SME_{source}_STAT" AS a WHERE report_date BETWEEN '{since}' AND '{until}' GROUP BY a.chartmetric_id ) AS b ON target.chartmetric_id = b.chartmetric_id AND target.run_date = CURRENT_DATE () WHEN matched THEN UPDATE SET target.{source}_max = b.max_date , target.{source} = b.processed_date WHEN NOT matched THEN INSERT ( chartmetric_id , run_date , status , {source}_max , {source}, schedule ) VALUES ( b.chartmetric_id , CURRENT_DATE () , 'BACKFILL' , b.max_date , b.processed_date, {int(schedule_value) + 1} ); """ return update_status_table_query def notify_slack(self): blocks = [{ "type": "section", "fields": [{ "type": "mrkdwn", "text": f"*Alert:*\n*Chartmetric Endpoint alert*\n" }, { "type": "mrkdwn", "text": f"*Type:*\n*Error*\n" }], "accessory": { "type": "image", "image_url": "https://api.slack.com/img/blocks/bkb_template_images" "/notificationsWarningIcon.png", "alt_text": "notifications warning icon" } }, { "type": "divider" }, { "type": "section", "fields": [{ "type": "mrkdwn", "text": f"\n*Environment:*\n{ENV.upper()}\n" }, { "type": "mrkdwn", "text": f"\n*Program End time:*\n{datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S')}\n" }], }, { "type": "divider" }, { "type": "section", "fields": [{ "type": "mrkdwn", "text": f"\n*Log Time*\n" }, { "type": "mrkdwn", "text": f"\n*Message*\n" }], }] if len(self.exception_records) > 0: for items in self.exception_records: for _key, _value in items.items(): blocks.append({ "type": "section", "fields": [{ "type": "mrkdwn", "text": f"{_key}\n" }, { "type": "mrkdwn", "text": f"{_value}\n" }] }) self.slack_client.chat_postMessage(channel=self.slack_channel, blocks=blocks) def log(self, message: str, _level: str): if _level == 'info': logger.info(message) elif _level == 'debug': logger.debug(message) elif _level == 'warning' or _level == 'error': logger.error(message) self.exception_records.append({ f"{datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S')}": str(message) }) if __name__ == "__main__": pass