Merged in dev_umistry (pull request #43)

Dev umistry
This commit is contained in:
Umang Mistry
2024-03-27 21:09:16 +00:00
5 changed files with 261 additions and 131 deletions
+123 -130
View File
@@ -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)
@@ -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)
@@ -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;
$$;
@@ -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
);
@@ -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.CLIENT_CONFIG_STAGE group by 1,2);
call stg.log_audit(:procedure_name, 'Section 4', 99, 'END');
RETURN 'Setup, Load, and Audit Complete';
END;
$$;