From d1d83f784207b61bc0123b9bf9c61d00336b95bd Mon Sep 17 00:00:00 2001 From: Umang Mistry Date: Mon, 25 Mar 2024 15:44:09 -0500 Subject: [PATCH 1/3] Added client config table DDL --- .../upload_interface/R__001_CLIENT_CONFIG_TABLE.sql | 7 +++++++ 1 file changed, 7 insertions(+) create mode 100644 snowflake/scripts/upload_interface/R__001_CLIENT_CONFIG_TABLE.sql diff --git a/snowflake/scripts/upload_interface/R__001_CLIENT_CONFIG_TABLE.sql b/snowflake/scripts/upload_interface/R__001_CLIENT_CONFIG_TABLE.sql new file mode 100644 index 0000000..42b03e7 --- /dev/null +++ b/snowflake/scripts/upload_interface/R__001_CLIENT_CONFIG_TABLE.sql @@ -0,0 +1,7 @@ +-- Create the table to store the client configuration details for Upload UI & Config UI +CREATE TABLE IF NOT EXISTS STG.CLIENT_CONFIG( + OA_CLIENT_ID NUMERIC, + CLIENT_NAME VARCHAR, + ACTIVE_PROJECT_COUNT NUMERIC, + S3_BUCKET_PATH VARCHAR +); \ No newline at end of file From cc59a1ccc12b33efe02e0b8f79558d221a8c7a81 Mon Sep 17 00:00:00 2001 From: Umang Mistry Date: Wed, 27 Mar 2024 16:06:06 -0500 Subject: [PATCH 2/3] Added business config objs & client config --- lambda/textract_logs_lambda.py | 253 +++++++++--------- .../training_data_pre_processing.py | 4 +- .../R__007_CREATE_BUSINESS_CONFIG_objects.sql | 70 +++++ .../R__002_LOAD_CLIENT_CONFIG_SP.sql | 58 ++++ 4 files changed, 254 insertions(+), 131 deletions(-) create mode 100644 snowflake/scripts/training_interface/R__007_CREATE_BUSINESS_CONFIG_objects.sql create mode 100644 snowflake/scripts/upload_interface/R__002_LOAD_CLIENT_CONFIG_SP.sql diff --git a/lambda/textract_logs_lambda.py b/lambda/textract_logs_lambda.py index 262edd1..0faa05d 100644 --- a/lambda/textract_logs_lambda.py +++ b/lambda/textract_logs_lambda.py @@ -1,54 +1,142 @@ import json import snowflake.connector +import boto3 +import snowflake.connector +import logging """ -Sample input event for Insert operation: -{ +# Sample input events + +doc_input_event = { "operation": "insert", + "table": "DOCUMENT_LOGS", "data": { - "BATCH_ID": 12345, - "JOB_ID": "job_67890", - "STAGE": "Preprocessing", - "TEXTRACT_STATUS": "Pending", - "BUCKET_NAME": "my-bucket", - "FILE_NAME": "document.pdf", - "FILE_PATH": "/path/to/document.pdf", - "DOCUMENT_TYPE": "Type A", - "PAYER_SIGNED": False, - "PROVIDER_SIGNED": True, - "GROUP_ID": "group_123", - "CREATED_TIME": "2023-01-01 12:00:00", - "MODIFIED_TIME": "2023-01-01 12:00:00", - "CREATED_BY": "user_1", - "MODIFIED_BY": "user_1", + "BATCH_ID": 101, + "JOB_ID": "J123456", + "STAGE": "Processing", + "TEXTRACT_STATUS": "Success", + "BUCKET_NAME": "doc-bucket", + "FILE_NAME": "file1.pdf", + "FILE_PATH": "/documents/2023/", + "DOCUMENT_TYPE": "Report", + "PAYER_SIGNED": True, + "PROVIDER_SIGNED": False, + "GROUP_ID": "G100", + "CREATED_TIME": "2024-03-21 10:00:00", + "MODIFIED_TIME": "2024-03-21 10:00:00", + "CREATED_BY": "admin", + "MODIFIED_BY": "admin", "ORIGINAL_FILE_EXTENSION": "pdf", "NO_OF_PAGES": 10, - "FILE_SIZE": 204800 + "FILE_SIZE": 1048576 } } -Sample input event for Update operation: - { - "operation": "update", - "data": { - "DOCUMENT_ID": 1001, - "TEXTRACT_STATUS": "Completed", - "MODIFIED_TIME": "2023-01-02 13:00:00", - "MODIFIED_BY": "user_2", - "NO_OF_PAGES": 12, - "FILE_SIZE": 304800 - } - } +doc_update_event = { + "operation": "update", + "table": "DOCUMENT_LOGS", + "data": { + "DOCUMENT_ID": 1001, + "TEXTRACT_STATUS": "Failed", + "MODIFIED_TIME": "2024-03-22 15:00:00", + "MODIFIED_BY": "admin" + } +} + + +batch_insert_event = { + "operation": "insert", + "table": "BATCH_LOGS", + "data": { + "CLIENT_ID": "C200", + "EXECUTION_START_TIME": "2024-03-21 09:00:00", + "NO_OF_DOCUMENTS": 150, + "USER_NAME": "batch_processor" + } +} + +batch_update_event = { + "operation": "update", + "table": "BATCH_LOGS", + "data": { + "BATCH_ID": 102, + "NO_OF_DOCUMENTS": 155, + "USER_NAME": "updated_processor" + } +} + +client_insert_event = { + "operation": "insert", + "table": "CLIENT_LOGS", + "data": { + "CLIENT_ID": "CL300", + "CLIENT_NAME": "Acme Corporation", + "BUCKET_NAME": "acme-docs" + } +} + +client_update_event = { + "operation": "update", + "table": "CLIENT_LOGS", + "data": { + "CLIENT_ID": "CL300", + "BUCKET_NAME": "new-acme-docs" + } +} Operation can be: 'insert' or 'update' Data is a dictionary with the columns and values to be inserted or updated """ +logging.basicConfig(level=logging.INFO, format='%(levelname)s: %(message)s', force=True) +logging.getLogger('snowflake.connector').setLevel(logging.WARNING) +logging.getLogger("botocore").setLevel(logging.WARNING) +logger = logging.getLogger(__name__) + + +def get_secret(secrets_name: str): + """Get credentials from Secret Manager as dict""" + secrets_manager = boto3.client('secretsmanager') + get_secret_value_response = secrets_manager.get_secret_value(SecretId=secrets_name) + + if 'SecretString' in get_secret_value_response: + secret_json = get_secret_value_response['SecretString'] + else: + secret_json = base64.b64decode(get_secret_value_response['SecretBinary']) + + return json.loads(secret_json) + + +def get_snowflake_db_connection(secrets_name: str): + """Create connection to Snowflake db.""" + try: + con_params = get_secret(secrets_name) + account = con_params['account_locator'] + user = con_params['user'] + password = con_params['password'] + database = con_params['database'].upper() + warehouse = con_params['warehouse'] + role= con_params['role'] + logger.info(f'Using credentials: account={account}, user={user}, password=***, database={database}, ' + f'warehouse={warehouse}') + snowflake_connection = snowflake.connector.connect(account=account, user=user, password=password, database=database, + warehouse=warehouse, autocommit=True) + logger.info(snowflake_connection) + return snowflake_connection + except Exception as e: + return e + + +# Leaving this statement outside the lambda_handler function to reuse the connection + # Secret has been setup to use the logging service account +conn = get_snowflake_db_connection('doczy-dev-db-svc-acc') +cur = conn.cursor() + def construct_doc_insert_sql(data): """ @@ -133,10 +221,8 @@ def lambda_handler(event, context): data = event['data'] table = event['table'] - # Get conn from Secrets Manager try: - # with conn.cursor() as cursor: if table == 'DOCUMENT_LOGS': if operation == 'insert': sql = construct_doc_insert_sql(data) @@ -159,108 +245,15 @@ def lambda_handler(event, context): sql = construct_client_update_sql(data, client_id) else: raise ValueError("Unsupported table.") - print(sql) - # cursor.execute(sql) + logger.info(f"Executing the following logging SQL statement: {sql}") + cur.execute(sql) return {'statusCode': 200, 'body': json.dumps('Operation successful')} except Exception as e: return {'statusCode': 400, 'body': json.dumps(str(e))} finally: # conn.close() + # Need to close the connection to avoid reaching the limit of open connections + # However, we need the conn to be in hot state for subsequent concurrent executions + # Need to decide the best approach to handle this pass - - -# Sample input events -doc_input_event = { - "operation": "insert", - "table": "DOCUMENT_LOGS", - "data": { - "BATCH_ID": 101, - "JOB_ID": "J123456", - "STAGE": "Processing", - "TEXTRACT_STATUS": "Success", - "BUCKET_NAME": "doc-bucket", - "FILE_NAME": "file1.pdf", - "FILE_PATH": "/documents/2023/", - "DOCUMENT_TYPE": "Report", - "PAYER_SIGNED": True, - "PROVIDER_SIGNED": False, - "GROUP_ID": "G100", - "CREATED_TIME": "2024-03-21 10:00:00", - "MODIFIED_TIME": "2024-03-21 10:00:00", - "CREATED_BY": "admin", - "MODIFIED_BY": "admin", - "ORIGINAL_FILE_EXTENSION": "pdf", - "NO_OF_PAGES": 10, - "FILE_SIZE": 1048576 - } -} - - - -doc_update_event = { - "operation": "update", - "table": "DOCUMENT_LOGS", - "data": { - "DOCUMENT_ID": 1001, - "TEXTRACT_STATUS": "Failed", - "MODIFIED_TIME": "2024-03-22 15:00:00", - "MODIFIED_BY": "admin" - } -} - - - -batch_insert_event = { - "operation": "insert", - "table": "BATCH_LOGS", - "data": { - "CLIENT_ID": "C200", - "EXECUTION_START_TIME": "2024-03-21 09:00:00", - "NO_OF_DOCUMENTS": 150, - "USER_NAME": "batch_processor" - } -} - -batch_update_event = { - "operation": "update", - "table": "BATCH_LOGS", - "data": { - "BATCH_ID": 102, - "NO_OF_DOCUMENTS": 155, - "USER_NAME": "updated_processor" - } -} - -client_insert_event = { - "operation": "insert", - "table": "CLIENT_LOGS", - "data": { - "CLIENT_ID": "CL300", - "CLIENT_NAME": "Acme Corporation", - "BUCKET_NAME": "acme-docs" - } -} - -client_update_event = { - "operation": "update", - "table": "CLIENT_LOGS", - "data": { - "CLIENT_ID": "CL300", - "BUCKET_NAME": "new-acme-docs" - } -} - - -lambda_handler(doc_input_event, None) - -lambda_handler(doc_update_event, None) - -lambda_handler(batch_insert_event, None) - -lambda_handler(batch_update_event, None) - -lambda_handler(client_insert_event, None) - -lambda_handler(client_update_event, None) - diff --git a/on_demand_scripts/training_data_pre_processing.py b/on_demand_scripts/training_data_pre_processing.py index 957f4fb..df9ac37 100644 --- a/on_demand_scripts/training_data_pre_processing.py +++ b/on_demand_scripts/training_data_pre_processing.py @@ -74,13 +74,14 @@ def create_business_config_table(file_name: str): xl_df = xl_df.iloc[:,12:] # Drop columns before DOCUMENT_NAME questions = xl_df.columns.tolist() # Grab the questions that are in the header row + field_name = xl_df.iloc[0].tolist() # Grab the field_name sf_cols = xl_df.iloc[1].tolist() # Grab the sf_cols priority = xl_df.iloc[3].tolist() # Grab the priority group_no = xl_df.iloc[4].tolist() # Grab the group_no theme = xl_df.iloc[5].tolist() # Grab the theme # Create a dataframe from the lists - df_internal = pd.DataFrame({'Column_name': sf_cols, 'Question': questions, 'priority': priority, 'group_no': group_no, 'theme': theme}) + df_internal = pd.DataFrame({'field_name':field_name,'Column_name': sf_cols, 'Question': questions, 'priority': priority, 'group_no': group_no, 'theme': theme}) # Drop rows where the question is 'Unnamed' and the column_name is NaN (Pandas automatically fills NaN with 'Unnamed' when reading excel files depending on the formatting) df_cleaned = df_internal[~df_internal['Question'].str.contains('Unnamed', na=False) & ~df_internal['Column_name'].isna()] @@ -89,6 +90,7 @@ def create_business_config_table(file_name: str): timestamp = str(date.strftime("%m%d%Y_%H%M%S")) df_cleaned.to_csv(f'biz_config-{timestamp}.csv', index=False) + return "Business config created successfully" except Exception as e: return str(e) diff --git a/snowflake/scripts/training_interface/R__007_CREATE_BUSINESS_CONFIG_objects.sql b/snowflake/scripts/training_interface/R__007_CREATE_BUSINESS_CONFIG_objects.sql new file mode 100644 index 0000000..e344076 --- /dev/null +++ b/snowflake/scripts/training_interface/R__007_CREATE_BUSINESS_CONFIG_objects.sql @@ -0,0 +1,70 @@ +CREATE TABLE IF NOT EXISTS STG.BUSINESS_CONFIG ( + FIELD_NAME VARCHAR + SF_COL_NAME VARCHAR, + QUESTION VARCHAR, + PRIORITY VARCHAR, + GROUP_NO VARCHAR, + THEME VARCHAR +); + + +CREATE OR REPLACE PROCEDURE STG.LOAD_BUSINESS_CONFIG(file_name VARCHAR) +RETURNS STRING +LANGUAGE SQL +EXECUTE AS CALLER +AS +$$ +DECLARE + procedure_name varchar; +BEGIN + + procedure_name := 'LOAD_BUSINESS_CONFIG'; + + call stg.log_audit(:procedure_name, 'Section 1', 99, 'START'); + + -- Create or replace stage with dynamic file name + EXECUTE IMMEDIATE 'CREATE OR REPLACE STAGE STG.BUSINESS_CONFIG_STAGE + STORAGE_INTEGRATION = dev_bucket_integration + URL = ''s3://doczy-dev-infra-raw-data-ingestion/training_data_raw/' || :file_name || ''' + FILE_FORMAT = (FORMAT_NAME = ''STG.CSV_HEADER'');'; + + call stg.log_audit(:procedure_name, 'Section 1', 99, 'END'); + + call stg.log_audit(:procedure_name, 'Section 2', 99, 'START'); + + -- Truncate table as we are using the KILL & FILL approach + TRUNCATE TABLE STG.BUSINESS_CONFIG; + + call stg.log_audit(:procedure_name, 'Section 2', 99, 'END'); + + call stg.log_audit(:procedure_name, 'Section 3', 99, 'START'); + + -- Copy command to load data + COPY INTO STG.BUSINESS_CONFIG FROM ( + SELECT + NULLIF(TRIM($1), '') AS FIELD_NAME, + NULLIF(TRIM($2), '') AS SF_COL_NAME, + NULLIF(TRIM($3), '') AS QUESTION, + NULLIF(TRIM($4), '') AS PRIORITY, + NULLIF(TRIM($5), '') AS GROUP_NO, + NULLIF(TRIM($6), '') AS THEME + FROM @STG.BUSINESS_CONFIG_STAGE + ) + FILE_FORMAT = (FORMAT_NAME = 'STG.CSV_HEADER') + ON_ERROR = ABORT_STATEMENT; + + call stg.log_audit(:procedure_name, 'Section 3', 99, 'END'); + + call stg.log_audit(:procedure_name, 'Section 4', 99, 'START'); + + INSERT INTO DOCZY_DEV.STG.DIM_AUDIT (AUDIT_SID, TABLE_NAME, SOURCE_FILE_NAME, LOAD_DATE, SOURCE_COUNT) + SELECT STG.AUDIT_SID.NEXTVAL, :procedure_name,* + FROM + (SELECT DISTINCT METADATA$FILENAME, CURRENT_TIMESTAMP(), max(METADATA$FILE_ROW_NUMBER) from @STG.BUSINESS_CONFIG_STAGE group by 1,2); + + call stg.log_audit(:procedure_name, 'Section 4', 99, 'END'); + + RETURN 'Setup, Load, and Audit Complete'; +END; +$$; + diff --git a/snowflake/scripts/upload_interface/R__002_LOAD_CLIENT_CONFIG_SP.sql b/snowflake/scripts/upload_interface/R__002_LOAD_CLIENT_CONFIG_SP.sql new file mode 100644 index 0000000..15d827b --- /dev/null +++ b/snowflake/scripts/upload_interface/R__002_LOAD_CLIENT_CONFIG_SP.sql @@ -0,0 +1,58 @@ +CREATE OR REPLACE PROCEDURE STG.LOAD_CLIENT_CONFIG(file_name VARCHAR) +RETURNS STRING +LANGUAGE SQL +EXECUTE AS CALLER +AS +$$ +DECLARE + procedure_name varchar; +BEGIN + + procedure_name := 'LOAD_CLIENT_CONFIG'; + + call stg.log_audit(:procedure_name, 'Section 1', 99, 'START'); + + -- Create or replace stage with dynamic file name + EXECUTE IMMEDIATE 'CREATE OR REPLACE STAGE STG.CLIENT_CONFIG_STAGE + STORAGE_INTEGRATION = dev_bucket_integration + URL = ''s3://doczy-dev-infra-raw-data-ingestion/client_names_openair/' || :file_name || ''' + FILE_FORMAT = (FORMAT_NAME = ''STG.CSV_HEADER'');'; + + call stg.log_audit(:procedure_name, 'Section 1', 99, 'END'); + + call stg.log_audit(:procedure_name, 'Section 2', 99, 'START'); + + -- Truncate table as we are using the KILL & FILL approach + TRUNCATE TABLE STG.CLIENT_CONFIG; + + call stg.log_audit(:procedure_name, 'Section 2', 99, 'END'); + + call stg.log_audit(:procedure_name, 'Section 3', 99, 'START'); + + -- Copy command to load data + COPY INTO STG.TRAINING_DATA_COLUMN_CONFIG FROM ( + SELECT + NULLIF(TRIM($1), '') AS OA_CLIENT_ID, + NULLIF(TRIM($2), '') AS CLIENT_NAME, + NULLIF(TRIM($3), '') AS ACTIVE_PROJECT_COUNT, + NULLIF(TRIM($4), '') AS S3_BUCKET_PATH, + FROM @STG.CLIENT_CONFIG_STAGE + ) + FILE_FORMAT = (FORMAT_NAME = 'STG.CSV_HEADER') + ON_ERROR = ABORT_STATEMENT; + + call stg.log_audit(:procedure_name, 'Section 3', 99, 'END'); + + call stg.log_audit(:procedure_name, 'Section 4', 99, 'START'); + + INSERT INTO DOCZY_DEV.STG.DIM_AUDIT (AUDIT_SID, TABLE_NAME, SOURCE_FILE_NAME, LOAD_DATE, SOURCE_COUNT) + SELECT STG.AUDIT_SID.NEXTVAL, :procedure_name,* + FROM + (SELECT DISTINCT METADATA$FILENAME, CURRENT_TIMESTAMP(), max(METADATA$FILE_ROW_NUMBER) from @STG.CSV_HEADER group by 1,2); + + call stg.log_audit(:procedure_name, 'Section 4', 99, 'END'); + + RETURN 'Setup, Load, and Audit Complete'; +END; +$$; + From 8d93106987b3df6eca88cda2fb26416b5a1aedd1 Mon Sep 17 00:00:00 2001 From: Umang Mistry Date: Wed, 27 Mar 2024 16:08:43 -0500 Subject: [PATCH 3/3] Updated typo in client config sp --- .../scripts/upload_interface/R__002_LOAD_CLIENT_CONFIG_SP.sql | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/snowflake/scripts/upload_interface/R__002_LOAD_CLIENT_CONFIG_SP.sql b/snowflake/scripts/upload_interface/R__002_LOAD_CLIENT_CONFIG_SP.sql index 15d827b..5744f42 100644 --- a/snowflake/scripts/upload_interface/R__002_LOAD_CLIENT_CONFIG_SP.sql +++ b/snowflake/scripts/upload_interface/R__002_LOAD_CLIENT_CONFIG_SP.sql @@ -48,7 +48,7 @@ BEGIN INSERT INTO DOCZY_DEV.STG.DIM_AUDIT (AUDIT_SID, TABLE_NAME, SOURCE_FILE_NAME, LOAD_DATE, SOURCE_COUNT) SELECT STG.AUDIT_SID.NEXTVAL, :procedure_name,* FROM - (SELECT DISTINCT METADATA$FILENAME, CURRENT_TIMESTAMP(), max(METADATA$FILE_ROW_NUMBER) from @STG.CSV_HEADER group by 1,2); + (SELECT DISTINCT METADATA$FILENAME, CURRENT_TIMESTAMP(), max(METADATA$FILE_ROW_NUMBER) from @STG.CLIENT_CONFIG_STAGE group by 1,2); call stg.log_audit(:procedure_name, 'Section 4', 99, 'END');