From 4f248040da3013f1f959d742440b8b10400bea5e Mon Sep 17 00:00:00 2001 From: Umang Mistry Date: Fri, 28 Jun 2024 17:13:28 -0500 Subject: [PATCH 1/3] Adding PROD tables for schemachange --- snowflake/PROD/R__2001_LOG_AUDIT.sql | 29 ++++++++ snowflake/PROD/R__2002_LOG_AUDIT_TABLE.sql | 11 +++ .../R__2001_UAT_PROMPT_CONFIG_TABLE.sql | 14 ++++ .../R__2002_UAT_CONTRACT_CONFIG_TABLE.sql | 16 +++++ .../R__2003_REQUEST_SUBMISSION_TABLE.sql | 6 ++ .../R__2004_LOAD_CONTRACT_CONFIG_SP.sql | 72 +++++++++++++++++++ .../R__2005_LOAD_REQUEST_SUBMISSION_SP.sql | 57 +++++++++++++++ .../R__2006_PIPELINE_OUTPUT.sql | 10 +++ ..._2007_LOAD_DOCZY_PIPELINE_RAW_OPUTPUTS.sql | 56 +++++++++++++++ .../R__2001_SERVERLESS_LOGS_TABLE.sql | 44 ++++++++++++ .../training_interface/R__2001_INIT_OBJS.sql | 66 +++++++++++++++++ .../R__2002_LOAD_TRAINING_RESULTS_SP.SQL | 58 +++++++++++++++ .../R__2003_LOAD_ATTEMPTS_SP.sql | 55 ++++++++++++++ .../R__2004_CREATE_TRAINING_DATA_TABLE_SP.sql | 28 ++++++++ .../R__2005_LOAD_COLUMN_CONFIG_SP.sql | 57 +++++++++++++++ .../R__2006_LOAD_TRAINING_DATA_RAW_SP.sql | 52 ++++++++++++++ ...R__2007_CREATE_BUSINESS_CONFIG_objects.sql | 8 +++ .../R__2008_LOAD_BUSINESS_CONFIG_SP.sql | 62 ++++++++++++++++ .../R__2001_CLIENT_CONFIG_TABLE.sql | 7 ++ .../R__2002_LOAD_CLIENT_CONFIG_SP.sql | 59 +++++++++++++++ .../R__2003_CONTRACT_UPLOAD_LOGS.sql | 10 +++ 21 files changed, 777 insertions(+) create mode 100644 snowflake/PROD/R__2001_LOG_AUDIT.sql create mode 100644 snowflake/PROD/R__2002_LOG_AUDIT_TABLE.sql create mode 100644 snowflake/PROD/config_interface/R__2001_UAT_PROMPT_CONFIG_TABLE.sql create mode 100644 snowflake/PROD/config_interface/R__2002_UAT_CONTRACT_CONFIG_TABLE.sql create mode 100644 snowflake/PROD/config_interface/R__2003_REQUEST_SUBMISSION_TABLE.sql create mode 100644 snowflake/PROD/config_interface/R__2004_LOAD_CONTRACT_CONFIG_SP.sql create mode 100644 snowflake/PROD/config_interface/R__2005_LOAD_REQUEST_SUBMISSION_SP.sql create mode 100644 snowflake/PROD/config_interface/R__2006_PIPELINE_OUTPUT.sql create mode 100644 snowflake/PROD/config_interface/R__2007_LOAD_DOCZY_PIPELINE_RAW_OPUTPUTS.sql create mode 100644 snowflake/PROD/system_wide/R__2001_SERVERLESS_LOGS_TABLE.sql create mode 100644 snowflake/PROD/training_interface/R__2001_INIT_OBJS.sql create mode 100644 snowflake/PROD/training_interface/R__2002_LOAD_TRAINING_RESULTS_SP.SQL create mode 100644 snowflake/PROD/training_interface/R__2003_LOAD_ATTEMPTS_SP.sql create mode 100644 snowflake/PROD/training_interface/R__2004_CREATE_TRAINING_DATA_TABLE_SP.sql create mode 100644 snowflake/PROD/training_interface/R__2005_LOAD_COLUMN_CONFIG_SP.sql create mode 100644 snowflake/PROD/training_interface/R__2006_LOAD_TRAINING_DATA_RAW_SP.sql create mode 100644 snowflake/PROD/training_interface/R__2007_CREATE_BUSINESS_CONFIG_objects.sql create mode 100644 snowflake/PROD/training_interface/R__2008_LOAD_BUSINESS_CONFIG_SP.sql create mode 100644 snowflake/PROD/upload_interface/R__2001_CLIENT_CONFIG_TABLE.sql create mode 100644 snowflake/PROD/upload_interface/R__2002_LOAD_CLIENT_CONFIG_SP.sql create mode 100644 snowflake/PROD/upload_interface/R__2003_CONTRACT_UPLOAD_LOGS.sql diff --git a/snowflake/PROD/R__2001_LOG_AUDIT.sql b/snowflake/PROD/R__2001_LOG_AUDIT.sql new file mode 100644 index 0000000..75fc24c --- /dev/null +++ b/snowflake/PROD/R__2001_LOG_AUDIT.sql @@ -0,0 +1,29 @@ +CREATE OR REPLACE PROCEDURE STG.LOG_AUDIT("PROC_NAME" VARCHAR(255), "SUB_SECTION_NAME" VARCHAR(255), "AUDIT_SID" FLOAT, "ACTION_FLAG" VARCHAR(255)) +RETURNS VARCHAR(16777216) +LANGUAGE SQL +EXECUTE AS CALLER +AS +' +BEGIN + BEGIN TRANSACTION; + + IF (:ACTION_FLAG = ''START'') THEN + INSERT INTO STG.LOG_AUDIT (EXECUTION_DT, START_DT, PROCEDURE_NAME, SUB_SECTION_NAME, AUDIT_SID) + VALUES (CURRENT_DATE(), CURRENT_TIMESTAMP(), :PROC_NAME, :SUB_SECTION_NAME, :AUDIT_SID); + END IF; + + IF (:ACTION_FLAG = ''END'') THEN + UPDATE STG.LOG_AUDIT + SET QUERY_ID = LAST_QUERY_ID(-2), END_DT = CURRENT_TIMESTAMP(), TOTAL_MIN = DATEDIFF(Minute, START_DT, CURRENT_TIMESTAMP()) + WHERE PROCEDURE_NAME = :PROC_NAME AND SUB_SECTION_NAME = :SUB_SECTION_NAME AND AUDIT_SID = :AUDIT_SID AND END_DT IS NULL; + + ELSE + COMMIT; + RETURN ''Incorrect Action Flag Used''; + + END IF; + + COMMIT; + RETURN ''Completed''; +END;' +; \ No newline at end of file diff --git a/snowflake/PROD/R__2002_LOG_AUDIT_TABLE.sql b/snowflake/PROD/R__2002_LOG_AUDIT_TABLE.sql new file mode 100644 index 0000000..d3ce540 --- /dev/null +++ b/snowflake/PROD/R__2002_LOG_AUDIT_TABLE.sql @@ -0,0 +1,11 @@ +create or replace TABLE STG.LOG_AUDIT ( + UNIQUE_KEY NUMBER(38,0) NOT NULL autoincrement start 1 increment 1 order, + EXECUTION_DT DATE, + QUERY_ID VARCHAR(100) COLLATE 'en-ci', + START_DT TIMESTAMP_NTZ(9), + END_DT TIMESTAMP_NTZ(9), + PROCEDURE_NAME VARCHAR(255) COLLATE 'en-ci', + SUB_SECTION_NAME VARCHAR(255) COLLATE 'en-ci', + AUDIT_SID NUMBER(38,0), + TOTAL_MIN NUMBER(38,0) +); diff --git a/snowflake/PROD/config_interface/R__2001_UAT_PROMPT_CONFIG_TABLE.sql b/snowflake/PROD/config_interface/R__2001_UAT_PROMPT_CONFIG_TABLE.sql new file mode 100644 index 0000000..01cfcdc --- /dev/null +++ b/snowflake/PROD/config_interface/R__2001_UAT_PROMPT_CONFIG_TABLE.sql @@ -0,0 +1,14 @@ +-- Separation of single vs multi prompt using a flag in the table +CREATE TABLE IF NOT EXISTS STG.PROMPT_CONFIG ( + FIELD_NAME VARCHAR, + FIELD_DESC VARCHAR, + IS_REQUIRED BOOLEAN, + FIELD_DATA_TYPE VARCHAR, + PROMPT VARCHAR, + FM_MODEL_ID VARCHAR, + GROUP_ID NUMERIC, + DEPENDENT_FIELDS VARCHAR, + FIELD_MAX_LENGTH NUMBER, + SAMPLE_VALUE VARCHAR, + CONTRACT_SECTION VARCHAR +); diff --git a/snowflake/PROD/config_interface/R__2002_UAT_CONTRACT_CONFIG_TABLE.sql b/snowflake/PROD/config_interface/R__2002_UAT_CONTRACT_CONFIG_TABLE.sql new file mode 100644 index 0000000..4770866 --- /dev/null +++ b/snowflake/PROD/config_interface/R__2002_UAT_CONTRACT_CONFIG_TABLE.sql @@ -0,0 +1,16 @@ +-- CREATING THE CONFIG TABLE +CREATE TABLE IF NOT EXISTS STG.CONTRACT_CONFIG ( + CONTRACT_NAME VARCHAR, + BATCH_ID VARCHAR, + UNIQUE_KEY BOOLEAN, + PRICING_BEFORE_CARVEOUTS BOOLEAN, + CONTRACT_RELATED BOOLEAN, + PROVIDER BOOLEAN, + TIMELINE BOOLEAN, + CARVEOUT_INDICATOR BOOLEAN, + CARVEOUT_METHODOLOGY BOOLEAN, + REQUEST_USER VARCHAR, + LATEST_FLAG BOOLEAN DEFAULT TRUE, + PIPELINE_KICKOFF_DATETIME DATETIME, + REQUEST_DATETIME DATETIME DEFAULT CURRENT_TIMESTAMP() +); \ No newline at end of file diff --git a/snowflake/PROD/config_interface/R__2003_REQUEST_SUBMISSION_TABLE.sql b/snowflake/PROD/config_interface/R__2003_REQUEST_SUBMISSION_TABLE.sql new file mode 100644 index 0000000..2107d73 --- /dev/null +++ b/snowflake/PROD/config_interface/R__2003_REQUEST_SUBMISSION_TABLE.sql @@ -0,0 +1,6 @@ +CREATE TABLE IF NOT EXISTS STG.REQUEST_SUBMISSION ( + CLIENT_NAME VARCHAR, + BATCH_ID VARCHAR, + REQUEST_USERNAME VARCHAR, + REQUEST_DATETIME DATETIME +); diff --git a/snowflake/PROD/config_interface/R__2004_LOAD_CONTRACT_CONFIG_SP.sql b/snowflake/PROD/config_interface/R__2004_LOAD_CONTRACT_CONFIG_SP.sql new file mode 100644 index 0000000..70aa801 --- /dev/null +++ b/snowflake/PROD/config_interface/R__2004_LOAD_CONTRACT_CONFIG_SP.sql @@ -0,0 +1,72 @@ +CREATE OR REPLACE PROCEDURE STG.LOAD_CONTRACT_CONFIG(file_name VARCHAR) +RETURNS STRING +LANGUAGE SQL +EXECUTE AS CALLER +AS +$$ +DECLARE + procedure_name varchar; +BEGIN + + procedure_name := 'LOAD_CONTRACT_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.CONTRACT_CONFIG_STAGE + STORAGE_INTEGRATION = prod_bucket_integration + URL = ''s3://doczyai-use2-p-infra-s3-raw-data-ingestion/doczy_pipeline_output/' || :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'); + COPY INTO STG.CONTRACT_CONFIG FROM ( + SELECT + NULLIF(TRIM($1),'') AS CONTRACT_NAME, + NULLIF(TRIM($2),'') AS BATCH_ID, + NULLIF(TRIM($3),'') AS UNIQUE_KEY, + NULLIF(TRIM($4),'') AS PRICING_BEFORE_CARVEOUTS, + NULLIF(TRIM($5),'') AS CONTRACT_RELATED, + NULLIF(TRIM($6),'') AS PROVIDER, + NULLIF(TRIM($7),'') AS TIMELINE, + NULLIF(TRIM($8),'') AS CARVEOUT_INDICATOR, + NULLIF(TRIM($9),'') AS CARVEOUT_METHODOLOGY, + NULLIF(TRIM($10),'') AS REQUEST_USER, + NULLIF(TRIM($11),'') AS LATEST_FLAG, + NULLIF(TRIM($12),'') AS PIPELINE_KICKOFF_DATETIME, + NULLIF(TRIM($13),'') AS REQUEST_DATETIME + FROM @STG.CONTRACT_CONFIG_STAGE + ) + FILE_FORMAT = (FORMAT_NAME = 'STG.CSV_HEADER') + ON_ERROR = ABORT_STATEMENT; + + call stg.log_audit(:procedure_name, 'Section 2', 99, 'END'); + + call stg.log_audit(:procedure_name, 'Section 3', 99, 'START'); + + -- Once the new data is loaded to the stage, we need to update the existing records in the table to set the LATEST_FLAG to FALSE + + contract_name := (SELECT NULLIF(TRIM($1),'') FROM @STG.CONTRACT_CONFIG_STAGE LIMIT 1); + batch_id := (SELECT NULLIF(TRIM($2),'') FROM @STG.CONTRACT_CONFIG_STAGE LIMIT 1); + request_datetime := (SELECT NULLIF(TRIM($13),'') FROM @STG.CONTRACT_CONFIG_STAGE LIMIT 1); + + UPDATE STG.CONTRACT_CONFIG + SET LATEST_FLAG = FALSE + WHERE BATCH_ID = :batch_id AND CONTRACT_NAME = :contract_name AND REQUEST_DATETIME < :request_datetime AND LATEST_FLAG = TRUE; + + call stg.log_audit(:procedure_name, 'Section 3', 99, 'END'); + + call stg.log_audit(:procedure_name, 'Section 4', 99, 'START'); + INSERT INTO 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.CONTRACT_CONFIG_STAGE group by 1,2); + + RETURN 'Setup, Load, and Audit Complete'; + + call stg.log_audit(:procedure_name, 'Section 4', 99, 'END'); + +END; +$$; + diff --git a/snowflake/PROD/config_interface/R__2005_LOAD_REQUEST_SUBMISSION_SP.sql b/snowflake/PROD/config_interface/R__2005_LOAD_REQUEST_SUBMISSION_SP.sql new file mode 100644 index 0000000..1254ae4 --- /dev/null +++ b/snowflake/PROD/config_interface/R__2005_LOAD_REQUEST_SUBMISSION_SP.sql @@ -0,0 +1,57 @@ +CREATE OR REPLACE PROCEDURE STG.LOAD_REQUEST_SUBMISSION(file_name VARCHAR) +RETURNS STRING +LANGUAGE SQL +EXECUTE AS CALLER +AS +$$ +DECLARE + procedure_name varchar; +BEGIN + + procedure_name := 'LOAD_REQUEST_SUBMISSION'; + + 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.REQUEST_SUBMISSION_STAGE + STORAGE_INTEGRATION = prod_bucket_integration + URL = ''s3://doczyai-use2-p-infra-s3-raw-data-ingestion/doczy_pipeline_output/' || :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'); + + COPY INTO STG.REQUEST_SUBMISSION FROM ( + SELECT + NULLIF(TRIM($2),'') AS CLIENT_NAME, + NULLIF(TRIM($3),'') AS BATCH_ID, + NULLIF(TRIM($4),'') AS REQUEST_USERNAME, + NULLIF(TRIM($5),'') AS REQUEST_DATETIME + + FROM @STG.REQUEST_SUBMISSION_STAGE + ) + FILE_FORMAT = (FORMAT_NAME = 'STG.CSV_HEADER') + ON_ERROR = ABORT_STATEMENT; + + call stg.log_audit(:procedure_name, 'Section 2', 99, 'END'); + + call stg.log_audit(:procedure_name, 'Section 3', 99, 'START'); + + INSERT INTO STG.DIM_AUDIT (AUDIT_SID, TABLE_NAME, SOURCE_FILE_NAME, LOAD_DATE, SOURCE_COUNT) + SELECT STG.AUDIT_SID.NEXTVAL, 'STG.REQUEST_SUBMISSION',* + FROM + (SELECT DISTINCT METADATA$FILENAME, CURRENT_TIMESTAMP(), max(METADATA$FILE_ROW_NUMBER) from @STG.REQUEST_SUBMISSION_STAGE group by 1,2); + + + + RETURN 'Setup, Load, and Audit Complete'; + + call stg.log_audit(:procedure_name, 'Section 3', 99, 'END'); + + + + +END; +$$; + diff --git a/snowflake/PROD/config_interface/R__2006_PIPELINE_OUTPUT.sql b/snowflake/PROD/config_interface/R__2006_PIPELINE_OUTPUT.sql new file mode 100644 index 0000000..b9849d1 --- /dev/null +++ b/snowflake/PROD/config_interface/R__2006_PIPELINE_OUTPUT.sql @@ -0,0 +1,10 @@ +-- CREATING THE CONFIG TABLE +CREATE TABLE IF NOT EXISTS STG.DOCZY_PIPELINE_RAW_OUTPUT_AC ( + CONTRACT_NAME VARCHAR(16777216) NOT NULL, + FIELD_NAME VARCHAR(255), + RAW_VALUE VARCHAR(16777216), + NEW_EXTRACTED_VALUE VARCHAR(16777216), + SNIPPET VARCHAR(16777216), + PAGE_NUMBER VARCHAR(255), + BATCH_ID VARCHAR(255) +); \ No newline at end of file diff --git a/snowflake/PROD/config_interface/R__2007_LOAD_DOCZY_PIPELINE_RAW_OPUTPUTS.sql b/snowflake/PROD/config_interface/R__2007_LOAD_DOCZY_PIPELINE_RAW_OPUTPUTS.sql new file mode 100644 index 0000000..8302a5d --- /dev/null +++ b/snowflake/PROD/config_interface/R__2007_LOAD_DOCZY_PIPELINE_RAW_OPUTPUTS.sql @@ -0,0 +1,56 @@ + CREATE OR REPLACE PROCEDURE STG.LOAD_DOCZY_PIPELINE_RAW_OUTPUT(file_name VARCHAR) +RETURNS STRING +LANGUAGE SQL +EXECUTE AS CALLER +AS +$$ +DECLARE + procedure_name VARCHAR; + unique_stage_name VARCHAR; +BEGIN + procedure_name := 'LOAD_DOCZY_PIPELINE_RAW_OUTPUT'; + unique_stage_name := 'DOCZY_PIPELINE_RAW_OUTPUT_STAGE_' || REPLACE(UUID_STRING(), '-', '_'); + + CALL stg.log_audit(:procedure_name, 'Section 1', 99, 'START'); + + -- Create a unique stage with dynamic file name + EXECUTE IMMEDIATE 'CREATE OR REPLACE STAGE STG.' || unique_stage_name || ' + STORAGE_INTEGRATION = PROD_BUCKET_INTEGRATION + URL = ''s3://doczyai-use2-p-infra-s3-raw-data-ingestion/doczy_pipeline_output/' || :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'); + + EXECUTE IMMEDIATE 'COPY INTO STG.DOCZY_PIPELINE_RAW_OUTPUT_AC FROM ( + SELECT + NULLIF(TRIM($1), '''') AS CONTRACT_NAME, + NULLIF(TRIM($2), '''') AS FIELD_NAME, + NULLIF(TRIM($3), '''') AS RAW_VALUE, + NULLIF(TRIM($4), '''') AS NEW_EXTRACTED_VALUE, + NULLIF(TRIM($5), '''') AS SNIPPET, + NULLIF(TRIM($6), '''') AS PAGE_NUMBER, + NULLIF(TRIM($7), '''') AS BATCH_ID + FROM @STG.' || unique_stage_name || ' + ) + FILE_FORMAT = (FORMAT_NAME = ''STG.CSV_HEADER'') + ON_ERROR = ABORT_STATEMENT;'; + + CALL stg.log_audit(:procedure_name, 'Section 2', 99, 'END'); + + CALL stg.log_audit(:procedure_name, 'Section 3', 99, 'START'); + + EXECUTE IMMEDIATE 'INSERT INTO STG.DIM_AUDIT (AUDIT_SID, TABLE_NAME, SOURCE_FILE_NAME, LOAD_DATE, SOURCE_COUNT) + SELECT STG.AUDIT_SID.NEXTVAL, ''STG.DOCZY_PIPELINE_RAW_OUTPUT_AC'', METADATA$FILENAME, CURRENT_TIMESTAMP(), MAX(METADATA$FILE_ROW_NUMBER) + FROM @STG.' || unique_stage_name || ' + GROUP BY 1, 2,3,4;'; + + -- Clean up the stage after use + EXECUTE IMMEDIATE 'LIST @STG.' || unique_stage_name; + + CALL stg.log_audit(:procedure_name, 'Section 3', 99, 'END'); + + RETURN 'Setup, Load, and Audit Complete'; +END; +$$; diff --git a/snowflake/PROD/system_wide/R__2001_SERVERLESS_LOGS_TABLE.sql b/snowflake/PROD/system_wide/R__2001_SERVERLESS_LOGS_TABLE.sql new file mode 100644 index 0000000..63b03cf --- /dev/null +++ b/snowflake/PROD/system_wide/R__2001_SERVERLESS_LOGS_TABLE.sql @@ -0,0 +1,44 @@ + + +-- Creating the 'document_logs' table +CREATE TABLE IF NOT EXISTS STG.DOCUMENT_LOGS +( + DOCUMENT_ID VARCHAR NOT NULL, + BATCH_ID VARCHAR NOT NULL, + JOB_ID VARCHAR, + STAGE VARCHAR, + TEXTRACT_STATUS VARCHAR, + BUCKET_NAME VARCHAR NOT NULL, + FILE_NAME VARCHAR NOT NULL, + FILE_PATH VARCHAR, + DOCUMENT_TYPE VARCHAR, + PAYER_SIGNED BOOLEAN, + PROVIDER_SIGNED BOOLEAN, + GROUP_ID VARCHAR, + CREATED_TIME TIMESTAMP, + MODIFIED_TIME TIMESTAMP, + CREATED_BY VARCHAR, + MODIFIED_BY VARCHAR, + ORIGINAL_FILE_EXTENSION VARCHAR, + NO_OF_PAGES NUMERIC, + FILE_SIZE NUMERIC, + PRIMARY KEY (DOCUMENT_ID) +); + +-- Creating the 'batch_logs' table +CREATE TABLE IF NOT EXISTS STG.BATCH_LOGS ( + BATCH_ID VARCHAR NOT NULL, + CLIENT_ID VARCHAR, + EXECUTION_START_TIME TIMESTAMP, + NO_OF_DOCUMENTS NUMBER, + USER_NAME VARCHAR, + PRIMARY KEY (BATCH_ID) +); + +-- Creating the 'client_logs' table +CREATE TABLE IF NOT EXISTS STG.CLIENT_LOGS ( + CLIENT_ID VARCHAR , + CLIENT_NAME VARCHAR, + BUCKET_NAME VARCHAR, + PRIMARY KEY (CLIENT_ID) +); diff --git a/snowflake/PROD/training_interface/R__2001_INIT_OBJS.sql b/snowflake/PROD/training_interface/R__2001_INIT_OBJS.sql new file mode 100644 index 0000000..992ece8 --- /dev/null +++ b/snowflake/PROD/training_interface/R__2001_INIT_OBJS.sql @@ -0,0 +1,66 @@ +-- Create staging table for results +CREATE TABLE IF NOT EXISTS STG.TRAINING_RESULTS +( + CONTRACT_NAME VARCHAR, + CONTRACT_ID NUMERIC, + ACTUAL_VALUE_STORED VARCHAR, + NEW_EXTRACTED_VALUE VARCHAR, + CONFIDENCE_LEVEL NUMERIC, + SNIPPET VARCHAR, + ORIGINAL_PAGE_NUMBER VARCHAR, + NEW_PAGE_NUMBER VARCHAR, + REVISED_PROMPT VARCHAR, + RESULT BOOLEAN, + LOAD_DT TIMESTAMP +); + + + -- Create file format +CREATE FILE FORMAT IF NOT EXISTS STG.CSV_HEADER + TYPE = 'CSV' + FIELD_DELIMITER = ',' + FILE_EXTENSION = '.csv' + RECORD_DELIMITER = '\\n' + DATE_FORMAT = AUTO + TRIM_SPACE = TRUE + NULL_IF = ('NULL', '', 'N/A','?','~','\\N') + SKIP_HEADER = 1 + EMPTY_FIELD_AS_NULL = TRUE + FIELD_OPTIONALLY_ENCLOSED_BY = '"' + SKIP_BLANK_LINES = TRUE; + + +-- Create staging table for attempt logs +CREATE TABLE IF NOT EXISTS STG.TRAINING_ATTEMPT_LOGS +( + FIELD_NAME VARCHAR, + CONTRACTS_TESTED NUMERIC, + USERNAME VARCHAR, + DATE_TIME TIMESTAMP, + ACCURACY NUMERIC, + ATTEMPT_NUM NUMERIC, + LOAD_DT TIMESTAMP +); + + + +-- Create table for column config that will be used by SP to create raw training data table +CREATE TABLE IF NOT EXISTS STG.TRAINING_DATA_COLUMN_CONFIG( +COLUMN_NAME VARCHAR, +COLUMN_DATATYPE VARCHAR +); + +-- Creating a separate file format for the training data +CREATE FILE FORMAT IF NOT EXISTS STG.TRAINING_DATA_FILE_FORMAT + TYPE = 'CSV' + FIELD_DELIMITER = ',' + FILE_EXTENSION = '.csv' + RECORD_DELIMITER = '\\n' + DATE_FORMAT = AUTO + TRIM_SPACE = TRUE + NULL_IF = ('NULL', '', 'N/A','?','~','\\N') + SKIP_HEADER = 1 + EMPTY_FIELD_AS_NULL = TRUE + FIELD_OPTIONALLY_ENCLOSED_BY = '"' + error_on_column_count_mismatch=false + SKIP_BLANK_LINES = TRUE; \ No newline at end of file diff --git a/snowflake/PROD/training_interface/R__2002_LOAD_TRAINING_RESULTS_SP.SQL b/snowflake/PROD/training_interface/R__2002_LOAD_TRAINING_RESULTS_SP.SQL new file mode 100644 index 0000000..af403c1 --- /dev/null +++ b/snowflake/PROD/training_interface/R__2002_LOAD_TRAINING_RESULTS_SP.SQL @@ -0,0 +1,58 @@ +CREATE OR REPLACE PROCEDURE STG.LOAD_TRAINING_RESULTS(file_name VARCHAR) +RETURNS STRING +LANGUAGE SQL +EXECUTE AS CALLER +AS +$$ +DECLARE + procedure_name varchar; +BEGIN + + procedure_name := 'LOAD_TRAINING_RESULTS'; + + call stg.log_audit(:procedure_name, 'Section 1', 99, 'START'); + + -- Create or replace stage with dynamic file name + -- TODO: Update for UAT + EXECUTE IMMEDIATE 'CREATE OR REPLACE STAGE STG.TRAINING_RESULTS_STAGE + STORAGE_INTEGRATION = prod_bucket_integration + URL = ''s3://doczyai-use2-p-infra-s3-raw-data-ingestion/doczy_pipeline_output/' || :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'); + -- Copy command to load data + COPY INTO STG.TRAINING_RESULTS FROM + ( + SELECT + NULLIF(TRIM($1), '') AS Contract_Name, + NULLIF(TRIM($2), '') AS Contract_ID, + NULLIF(TRIM($3), '') AS Actual_Value_Stored, + NULLIF(TRIM($4), '') AS New_Extracted_Value, + NULLIF(TRIM($5), '') AS Confidence_Level, + NULLIF(TRIM($6), '') AS Snippet, + NULLIF(TRIM($7), '') AS Original_Page_Number, + NULLIF(TRIM($8), '') AS New_Page_Number, + NULLIF(TRIM($9), '') AS Revised_Prompt, + IFF(TRIM($10) = 'TRUE', TRUE, FALSE) AS Result, + CURRENT_TIMESTAMP(2) AS LOAD_DT + FROM @STG.TRAINING_RESULTS_STAGE + ) + FILE_FORMAT = (FORMAT_NAME = 'STG.CSV_HEADER') + ON_ERROR = ABORT_STATEMENT; + + call stg.log_audit(:procedure_name, 'Section 2', 99, 'END'); + + + call stg.log_audit(:procedure_name, 'Section 3', 99, 'START'); + + INSERT INTO STG.DIM_AUDIT (AUDIT_SID, TABLE_NAME, SOURCE_FILE_NAME, LOAD_DATE, SOURCE_COUNT) + SELECT STG.AUDIT_SID.NEXTVAL, 'STG.TRAINING_RESULTS',* + FROM + (SELECT DISTINCT METADATA$FILENAME, CURRENT_TIMESTAMP(), max(METADATA$FILE_ROW_NUMBER) from @STG.TRAINING_RESULTS_STAGE group by 1,2); + + call stg.log_audit(:procedure_name, 'Section 3', 99, 'END'); + + RETURN 'Setup, Load, and Audit Complete'; +END; +$$; \ No newline at end of file diff --git a/snowflake/PROD/training_interface/R__2003_LOAD_ATTEMPTS_SP.sql b/snowflake/PROD/training_interface/R__2003_LOAD_ATTEMPTS_SP.sql new file mode 100644 index 0000000..441aa0e --- /dev/null +++ b/snowflake/PROD/training_interface/R__2003_LOAD_ATTEMPTS_SP.sql @@ -0,0 +1,55 @@ +CREATE OR REPLACE PROCEDURE STG.LOAD_ATTEMPT_LOGS(file_name VARCHAR) +RETURNS STRING +LANGUAGE SQL +EXECUTE AS CALLER +AS +$$ +DECLARE + procedure_name varchar; +BEGIN + + procedure_name := 'LOAD_ATTEMPT_LOGS'; + + call stg.log_audit(:procedure_name, 'Section 1', 99, 'START'); + -- Create or replace stage with dynamic file name + -- TODO: Update for UAT + EXECUTE IMMEDIATE 'CREATE OR REPLACE STAGE STG.ATTEMPT_LOGS_STAGE + STORAGE_INTEGRATION = prod_bucket_integration + URL = ''s3://doczyai-use2-p-infra-s3-raw-data-ingestion/doczy_pipeline_output/' || :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'); + + -- Copy command to load data + COPY INTO STG.TRAINING_ATTEMPT_LOGS FROM + ( + SELECT + NULLIF(TRIM($1), '') AS FIELD_NAME, + NULLIF(TRIM($2), '') AS CONTRACTS_TESTED, + NULLIF(TRIM($3), '') AS USERNAME, + NULLIF(TRIM($4), '') AS DATE_TIME, + NULLIF(TRIM($5), '') AS ACCURACY, + NULLIF(TRIM($6), '') AS ATTEMPT_NUM, + CURRENT_TIMESTAMP(2) AS LOAD_DT + FROM @STG.ATTEMPT_LOGS_STAGE + ) + FILE_FORMAT = (FORMAT_NAME = 'STG.CSV_HEADER') + ON_ERROR = ABORT_STATEMENT; + + call stg.log_audit(:procedure_name, 'Section 2', 99, 'END'); + + call stg.log_audit(:procedure_name, 'Section 3', 99, 'START'); + + INSERT INTO STG.DIM_AUDIT (AUDIT_SID, TABLE_NAME, SOURCE_FILE_NAME, LOAD_DATE, SOURCE_COUNT) + SELECT STG.AUDIT_SID.NEXTVAL, 'STG.TRAINING_ATTEMPT_LOGS',* + FROM + (SELECT DISTINCT METADATA$FILENAME, CURRENT_TIMESTAMP(), max(METADATA$FILE_ROW_NUMBER) from @STG.ATTEMPT_LOGS_STAGE group by 1,2); + + call stg.log_audit(:procedure_name, 'Section 3', 99, 'END'); + + RETURN 'Setup, Load, and Audit Complete'; +END; +$$; + diff --git a/snowflake/PROD/training_interface/R__2004_CREATE_TRAINING_DATA_TABLE_SP.sql b/snowflake/PROD/training_interface/R__2004_CREATE_TRAINING_DATA_TABLE_SP.sql new file mode 100644 index 0000000..eeb3dbc --- /dev/null +++ b/snowflake/PROD/training_interface/R__2004_CREATE_TRAINING_DATA_TABLE_SP.sql @@ -0,0 +1,28 @@ +CREATE OR REPLACE PROCEDURE STG.CREATE_TRAINING_DATA_TABLE() +RETURNS VARCHAR(16777216) +LANGUAGE SQL +EXECUTE AS CALLER +AS ' +DECLARE + dynamic_ddl STRING := ''CREATE OR REPLACE TABLE STG.TRAINING_DATA_RAW (''; + column_details RESULTSET; + first_column BOOLEAN := TRUE; + cur_config cursor FOR + SELECT column_name, column_datatype FROM STG.TRAINING_DATA_COLUMN_CONFIG; +BEGIN + -- Creating the raw training data table with all columns as VARCHAR due to the dynamic nature of the columns and fields + OPEN cur_config; + FOR rec IN cur_config DO + dynamic_ddl := dynamic_ddl || rec.column_name || '' '' || ''VARCHAR'' || '',''; + END FOR; + + dynamic_ddl := LEFT(dynamic_ddl, LENGTH(dynamic_ddl) - 1); + + dynamic_ddl := dynamic_ddl || '');''; + + EXECUTE IMMEDIATE dynamic_ddl; + + RETURN dynamic_ddl; +END; + +'; \ No newline at end of file diff --git a/snowflake/PROD/training_interface/R__2005_LOAD_COLUMN_CONFIG_SP.sql b/snowflake/PROD/training_interface/R__2005_LOAD_COLUMN_CONFIG_SP.sql new file mode 100644 index 0000000..8fef6db --- /dev/null +++ b/snowflake/PROD/training_interface/R__2005_LOAD_COLUMN_CONFIG_SP.sql @@ -0,0 +1,57 @@ +CREATE OR REPLACE PROCEDURE STG.LOAD_COLUMN_CONFIG(file_name VARCHAR) +RETURNS STRING +LANGUAGE SQL +EXECUTE AS CALLER +AS +$$ +DECLARE + procedure_name varchar; +BEGIN + + procedure_name := 'LOAD_COLUMN_CONFIG'; + + call stg.log_audit(:procedure_name, 'Section 1', 99, 'START'); + + -- Create or replace stage with dynamic file name + -- TODO: Update for UAT + EXECUTE IMMEDIATE 'CREATE OR REPLACE STAGE STG.COLUMN_CONFIG_STAGE + STORAGE_INTEGRATION = prod_bucket_integration + URL = ''s3://doczyai-use2-p-infra-s3-raw-data-ingestion/doczy_pipeline_output/' || :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.TRAINING_DATA_COLUMN_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 COLUMN_NAME, + NULLIF(TRIM($2), '') AS COLUMN_DATATYPE + FROM @STG.COLUMN_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 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.COLUMN_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/PROD/training_interface/R__2006_LOAD_TRAINING_DATA_RAW_SP.sql b/snowflake/PROD/training_interface/R__2006_LOAD_TRAINING_DATA_RAW_SP.sql new file mode 100644 index 0000000..fa3285c --- /dev/null +++ b/snowflake/PROD/training_interface/R__2006_LOAD_TRAINING_DATA_RAW_SP.sql @@ -0,0 +1,52 @@ +CREATE OR REPLACE PROCEDURE STG.LOAD_TRAINING_DATA_RAW(file_name VARCHAR) +RETURNS STRING +LANGUAGE SQL +EXECUTE AS CALLER +AS +$$ +DECLARE + procedure_name varchar; +BEGIN + + procedure_name := 'LOAD_TRAINING_DATA_RAW'; + + call stg.log_audit(:procedure_name, 'Section 1', 99, 'START'); + + -- Create or replace stage with dynamic file name + -- TODO: Update for UAT + EXECUTE IMMEDIATE 'CREATE OR REPLACE STAGE STG.RAW_TRAINING_DATA_STAGE + STORAGE_INTEGRATION = prod_bucket_integration + URL = ''s3://doczyai-use2-p-infra-s3-raw-data-ingestion/doczy_pipeline_output/' || :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'); + + -- Recreating the training data table by calling the SP. This will replace the existing table with new column definitions + call STG.CREATE_TRAINING_DATA_TABLE(); + + 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_RAW FROM @STG.RAW_TRAINING_DATA_STAGE + FILE_FORMAT = (FORMAT_NAME = 'STG.TRAINING_DATA_FILE_FORMAT') + 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 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.RAW_TRAINING_DATA_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/PROD/training_interface/R__2007_CREATE_BUSINESS_CONFIG_objects.sql b/snowflake/PROD/training_interface/R__2007_CREATE_BUSINESS_CONFIG_objects.sql new file mode 100644 index 0000000..2601408 --- /dev/null +++ b/snowflake/PROD/training_interface/R__2007_CREATE_BUSINESS_CONFIG_objects.sql @@ -0,0 +1,8 @@ +CREATE TABLE IF NOT EXISTS STG.BUSINESS_CONFIG ( + FIELD_NAME VARCHAR , + SF_COL_NAME VARCHAR, + QUESTION VARCHAR, + PRIORITY VARCHAR, + GROUP_NO VARCHAR, + THEME VARCHAR +); \ No newline at end of file diff --git a/snowflake/PROD/training_interface/R__2008_LOAD_BUSINESS_CONFIG_SP.sql b/snowflake/PROD/training_interface/R__2008_LOAD_BUSINESS_CONFIG_SP.sql new file mode 100644 index 0000000..83cff7a --- /dev/null +++ b/snowflake/PROD/training_interface/R__2008_LOAD_BUSINESS_CONFIG_SP.sql @@ -0,0 +1,62 @@ + +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 + -- TODO: Update for UAT + EXECUTE IMMEDIATE 'CREATE OR REPLACE STAGE STG.BUSINESS_CONFIG_STAGE + STORAGE_INTEGRATION = prod_bucket_integration + URL = ''s3://doczyai-use2-p-infra-s3-raw-data-ingestion/doczy_pipeline_output/' || :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 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/PROD/upload_interface/R__2001_CLIENT_CONFIG_TABLE.sql b/snowflake/PROD/upload_interface/R__2001_CLIENT_CONFIG_TABLE.sql new file mode 100644 index 0000000..42b03e7 --- /dev/null +++ b/snowflake/PROD/upload_interface/R__2001_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 diff --git a/snowflake/PROD/upload_interface/R__2002_LOAD_CLIENT_CONFIG_SP.sql b/snowflake/PROD/upload_interface/R__2002_LOAD_CLIENT_CONFIG_SP.sql new file mode 100644 index 0000000..4f15448 --- /dev/null +++ b/snowflake/PROD/upload_interface/R__2002_LOAD_CLIENT_CONFIG_SP.sql @@ -0,0 +1,59 @@ +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 + -- TODO: Update for UAT + EXECUTE IMMEDIATE 'CREATE OR REPLACE STAGE STG.CLIENT_CONFIG_STAGE + STORAGE_INTEGRATION = prod_bucket_integration + URL = ''s3://doczyai-use2-p-infra-s3-raw-data-ingestion/doczy_pipeline_output/' || :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.CLIENT_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 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; +$$; + diff --git a/snowflake/PROD/upload_interface/R__2003_CONTRACT_UPLOAD_LOGS.sql b/snowflake/PROD/upload_interface/R__2003_CONTRACT_UPLOAD_LOGS.sql new file mode 100644 index 0000000..ac7126e --- /dev/null +++ b/snowflake/PROD/upload_interface/R__2003_CONTRACT_UPLOAD_LOGS.sql @@ -0,0 +1,10 @@ +-- This table is used to store the logs of all the files that were uploaded by the user to S3 +-- Following fields are to be stored in the table: +-- Batch id & client name & file name, date time, user +CREATE TABLE IF NOT EXISTS STG.CONTRACT_UPLOAD_LOGS ( + BATCH_ID VARCHAR, + CLIENT_NAME VARCHAR, + FILE_NAME VARCHAR, + UPLOAD_DATETIME DATETIME DEFAULT CURRENT_TIMESTAMP(), + UPLOAD_USER VARCHAR +); \ No newline at end of file From efb38bece93712ad4d8314dcfca3795e9600632a Mon Sep 17 00:00:00 2001 From: Umang Mistry Date: Fri, 28 Jun 2024 17:38:24 -0500 Subject: [PATCH 2/3] Added B fields raw output ingestion --- .../R__2008_PIPELINE_OUTPUT_B_FIELDS.sql | 61 ++++++++++ ...2009_LOAD_B_FIELDS_PIPELINE_RAW_OUTPUT.sql | 108 ++++++++++++++++++ 2 files changed, 169 insertions(+) create mode 100644 snowflake/PROD/config_interface/R__2008_PIPELINE_OUTPUT_B_FIELDS.sql create mode 100644 snowflake/PROD/config_interface/R__2009_LOAD_B_FIELDS_PIPELINE_RAW_OUTPUT.sql diff --git a/snowflake/PROD/config_interface/R__2008_PIPELINE_OUTPUT_B_FIELDS.sql b/snowflake/PROD/config_interface/R__2008_PIPELINE_OUTPUT_B_FIELDS.sql new file mode 100644 index 0000000..d02885f --- /dev/null +++ b/snowflake/PROD/config_interface/R__2008_PIPELINE_OUTPUT_B_FIELDS.sql @@ -0,0 +1,61 @@ +CREATE TABLE IF NOT EXISTS STG.DOCZY_PIPELINE_RAW_OUTPUT_B ( + SERVICE TEXT, + REIMBURSEMENT_FLAT_FEE TEXT, + REIMBURSEMENT_RATE TEXT, + FULL_METHODOLOGY TEXT, + page_num TEXT, + contract_name TEXT, + LESSER_OF_LANGUAGE_IND TEXT, + GREATER_OF_LANGUAGE_IND TEXT, + REIMBURSEMENT_METHODOLOGY TEXT, + REIMBURSEMENT_FEE_SCHEDULE TEXT, + REIMBURSEMENT_FEE_SCHEDULE_VERSION TEXT, + REIMBURSEMENT_EXCEPTION_IND TEXT, + REIMBURSEMENT_DESCRIBE_EXCEPTION TEXT, + REIMBURSEMENT_ADMITTYPE_CODES TEXT, + REIMBURSEMENT_DIAG_CODES TEXT, + REIMBURSEMENT_GROUPER_CODES TEXT, + REIMBURSEMENT_GROUPER TEXT, + REIMBURSEMENT_PLACEOFSERVICE_CODES TEXT, + REIMBURSEMENT_PROC_CODES TEXT, + REIMBURSEMENT_REVENUE_CODES TEXT, + REIMBURSEMENT_STATUS_INDICATOR_CODES TEXT, + SERVICE_PG TEXT, + REIMBURSEMENT_FLAT_FEE_PG TEXT, + REIMBURSEMENT_RATE_PG TEXT, + FULL_METHODOLOGY_PG TEXT, + Filename_PG TEXT, + LESSER_OF_LANGUAGE_IND_PG TEXT, + GREATER_OF_LANGUAGE_IND_PG TEXT, + REIMBURSEMENT_METHODOLOGY_PG TEXT, + REIMBURSEMENT_FEE_SCHEDULE_PG TEXT, + REIMBURSEMENT_FEE_SCHEDULE_VERSION_PG TEXT, + REIMBURSEMENT_EXCEPTION_IND_PG TEXT, + REIMBURSEMENT_DESCRIBE_EXCEPTION_PG TEXT, + REIMBURSEMENT_ADMITTYPE_CODES_PG TEXT, + REIMBURSEMENT_DIAG_CODES_PG TEXT, + REIMBURSEMENT_GROUPER_CODES_PG TEXT, + REIMBURSEMENT_GROUPER_PG TEXT, + REIMBURSEMENT_PLACEOFSERVICE_CODES_PG TEXT, + REIMBURSEMENT_PROC_CODES_PG TEXT, + REIMBURSEMENT_REVENUE_CODES_PG TEXT, + REIMBURSEMENT_STATUS_INDICATOR_CODES_PG TEXT, + CONTRACT_LOB TEXT, + CONTRACT_LOB_PG TEXT, + CONTRACT_PROGRAM TEXT, + CONTRACT_PROGRAM_PG TEXT, + CONTRACT_NETWORK TEXT, + CONTRACT_NETWORK_PG TEXT, + PRODUCT TEXT, + PRODUCT_PG TEXT, + CONTRACT_MARKETPLACE_METAL_LEVEL TEXT, + CONTRACT_MARKETPLACE_METAL_LEVEL_PG TEXT, + LOB_PRICING_TERMS_EFFECTIVE_DATE TEXT, + LOB_PRICING_TERMS_EFFECTIVE_DATE_PG TEXT, + LOB_PRICING_TERMS_TERMINATION_DATE TEXT, + LOB_PRICING_TERMS_TERMINATION_DATE_PG TEXT, + Corrected_LOB TEXT, + Corrected_PROGRAM TEXT, + Corrected_NETWORK TEXT, + batch_id TEXT +); diff --git a/snowflake/PROD/config_interface/R__2009_LOAD_B_FIELDS_PIPELINE_RAW_OUTPUT.sql b/snowflake/PROD/config_interface/R__2009_LOAD_B_FIELDS_PIPELINE_RAW_OUTPUT.sql new file mode 100644 index 0000000..69139dc --- /dev/null +++ b/snowflake/PROD/config_interface/R__2009_LOAD_B_FIELDS_PIPELINE_RAW_OUTPUT.sql @@ -0,0 +1,108 @@ +CREATE OR REPLACE PROCEDURE STG.LOAD_DOCZY_PIPELINE_RAW_OUTPUT_B_FIELDS(file_name VARCHAR) +RETURNS STRING +LANGUAGE SQL +EXECUTE AS CALLER +AS +$$ +DECLARE + procedure_name VARCHAR; + unique_stage_name VARCHAR; +BEGIN + procedure_name := 'LOAD_DOCZY_PIPELINE_RAW_OUTPUT_LARGE'; + unique_stage_name := 'DOCZY_PIPELINE_RAW_OUTPUT_STAGE_' || REPLACE(UUID_STRING(), '-', '_'); + + CALL stg.log_audit(:procedure_name, 'Section 1', 99, 'START'); + + -- Create a unique stage with dynamic file name + EXECUTE IMMEDIATE 'CREATE OR REPLACE STAGE STG.' || unique_stage_name || ' + STORAGE_INTEGRATION = PROD_BUCKET_INTEGRATION + URL = ''s3://doczyai-use2-p-infra-s3-raw-data-ingestion/doczy_pipeline_output/' || :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'); + + EXECUTE IMMEDIATE 'COPY INTO DOCZY_PIPELINE_RAW_OUTPUT_B FROM ( + SELECT + NULLIF(TRIM($1), '''') AS SERVICE, + NULLIF(TRIM($2), '''') AS REIMBURSEMENT_FLAT_FEE, + NULLIF(TRIM($3), '''') AS REIMBURSEMENT_RATE, + NULLIF(TRIM($4), '''') AS FULL_METHODOLOGY, + NULLIF(TRIM($5), '''') AS page_num, + NULLIF(TRIM($6), '''') AS contract_name, + NULLIF(TRIM($7), '''') AS LESSER_OF_LANGUAGE_IND, + NULLIF(TRIM($8), '''') AS GREATER_OF_LANGUAGE_IND, + NULLIF(TRIM($9), '''') AS REIMBURSEMENT_METHODOLOGY, + NULLIF(TRIM($10), '''') AS REIMBURSEMENT_FEE_SCHEDULE, + NULLIF(TRIM($11), '''') AS REIMBURSEMENT_FEE_SCHEDULE_VERSION, + NULLIF(TRIM($12), '''') AS REIMBURSEMENT_EXCEPTION_IND, + NULLIF(TRIM($13), '''') AS REIMBURSEMENT_DESCRIBE_EXCEPTION, + NULLIF(TRIM($14), '''') AS REIMBURSEMENT_ADMITTYPE_CODES, + NULLIF(TRIM($15), '''') AS REIMBURSEMENT_DIAG_CODES, + NULLIF(TRIM($16), '''') AS REIMBURSEMENT_GROUPER_CODES, + NULLIF(TRIM($17), '''') AS REIMBURSEMENT_GROUPER, + NULLIF(TRIM($18), '''') AS REIMBURSEMENT_PLACEOFSERVICE_CODES, + NULLIF(TRIM($19), '''') AS REIMBURSEMENT_PROC_CODES, + NULLIF(TRIM($20), '''') AS REIMBURSEMENT_REVENUE_CODES, + NULLIF(TRIM($21), '''') AS REIMBURSEMENT_STATUS_INDICATOR_CODES, + NULLIF(TRIM($22), '''') AS SERVICE_PG, + NULLIF(TRIM($23), '''') AS REIMBURSEMENT_FLAT_FEE_PG, + NULLIF(TRIM($24), '''') AS REIMBURSEMENT_RATE_PG, + NULLIF(TRIM($25), '''') AS FULL_METHODOLOGY_PG, + NULLIF(TRIM($26), '''') AS contract_name_PG, + NULLIF(TRIM($27), '''') AS LESSER_OF_LANGUAGE_IND_PG, + NULLIF(TRIM($28), '''') AS GREATER_OF_LANGUAGE_IND_PG, + NULLIF(TRIM($29), '''') AS REIMBURSEMENT_METHODOLOGY_PG, + NULLIF(TRIM($30), '''') AS REIMBURSEMENT_FEE_SCHEDULE_PG, + NULLIF(TRIM($31), '''') AS REIMBURSEMENT_FEE_SCHEDULE_VERSION_PG, + NULLIF(TRIM($32), '''') AS REIMBURSEMENT_EXCEPTION_IND_PG, + NULLIF(TRIM($33), '''') AS REIMBURSEMENT_DESCRIBE_EXCEPTION_PG, + NULLIF(TRIM($34), '''') AS REIMBURSEMENT_ADMITTYPE_CODES_PG, + NULLIF(TRIM($35), '''') AS REIMBURSEMENT_DIAG_CODES_PG, + NULLIF(TRIM($36), '''') AS REIMBURSEMENT_GROUPER_CODES_PG, + NULLIF(TRIM($37), '''') AS REIMBURSEMENT_GROUPER_PG, + NULLIF(TRIM($38), '''') AS REIMBURSEMENT_PLACEOFSERVICE_CODES_PG, + NULLIF(TRIM($39), '''') AS REIMBURSEMENT_PROC_CODES_PG, + NULLIF(TRIM($40), '''') AS REIMBURSEMENT_REVENUE_CODES_PG, + NULLIF(TRIM($41), '''') AS REIMBURSEMENT_STATUS_INDICATOR_CODES_PG, + NULLIF(TRIM($42), '''') AS CONTRACT_LOB, + NULLIF(TRIM($43), '''') AS CONTRACT_LOB_PG, + NULLIF(TRIM($44), '''') AS CONTRACT_PROGRAM, + NULLIF(TRIM($45), '''') AS CONTRACT_PROGRAM_PG, + NULLIF(TRIM($46), '''') AS CONTRACT_NETWORK, + NULLIF(TRIM($47), '''') AS CONTRACT_NETWORK_PG, + NULLIF(TRIM($48), '''') AS PRODUCT, + NULLIF(TRIM($49), '''') AS PRODUCT_PG, + NULLIF(TRIM($50), '''') AS CONTRACT_MARKETPLACE_METAL_LEVEL, + NULLIF(TRIM($51), '''') AS CONTRACT_MARKETPLACE_METAL_LEVEL_PG, + NULLIF(TRIM($52), '''') AS LOB_PRICING_TERMS_EFFECTIVE_DATE, + NULLIF(TRIM($53), '''') AS LOB_PRICING_TERMS_EFFECTIVE_DATE_PG, + NULLIF(TRIM($54), '''') AS LOB_PRICING_TERMS_TERMINATION_DATE, + NULLIF(TRIM($55), '''') AS LOB_PRICING_TERMS_TERMINATION_DATE_PG, + NULLIF(TRIM($56), '''') AS Corrected_LOB, + NULLIF(TRIM($57), '''') AS Corrected_PROGRAM, + NULLIF(TRIM($58), '''') AS Corrected_NETWORK, + NULLIF(TRIM($59), '''') AS batch_id + FROM @STG.' || unique_stage_name || ' + ) + FILE_FORMAT = (FORMAT_NAME = ''STG.CSV_HEADER'') + ON_ERROR = ABORT_STATEMENT;'; + + CALL stg.log_audit(:procedure_name, 'Section 2', 99, 'END'); + + CALL stg.log_audit(:procedure_name, 'Section 3', 99, 'START'); + + EXECUTE IMMEDIATE 'INSERT INTO STG.DIM_AUDIT (AUDIT_SID, TABLE_NAME, SOURCE_FILE_NAME, LOAD_DATE, SOURCE_COUNT) + SELECT STG.AUDIT_SID.NEXTVAL, ''DOCZY_PIPELINE_RAW_OUTPUT_B'', METADATA$FILENAME, CURRENT_TIMESTAMP(), MAX(METADATA$FILE_ROW_NUMBER) + FROM @STG.' || unique_stage_name || ' + GROUP BY 1, 2,3,4;'; + + -- Clean up the stage after use + EXECUTE IMMEDIATE 'DROP STAGE IF EXISTS STG.' || unique_stage_name; + + CALL stg.log_audit(:procedure_name, 'Section 3', 99, 'END'); + + RETURN 'Setup, Load, and Audit Complete'; +END; +$$; From 9268a55184cab02f9b383e90e101e15d1638f54b Mon Sep 17 00:00:00 2001 From: Umang Mistry Date: Fri, 28 Jun 2024 17:45:28 -0500 Subject: [PATCH 3/3] Updated lambda handler & removed keys from config --- .../src/lambda/prompt-orchestrator/config.py | 14 +-- .../src/lambda/prompt-orchestrator/index.py | 112 ++++++++++++++++-- 2 files changed, 104 insertions(+), 22 deletions(-) diff --git a/textract-pipeline/src/lambda/prompt-orchestrator/config.py b/textract-pipeline/src/lambda/prompt-orchestrator/config.py index 02baf9f..0ee313b 100644 --- a/textract-pipeline/src/lambda/prompt-orchestrator/config.py +++ b/textract-pipeline/src/lambda/prompt-orchestrator/config.py @@ -75,30 +75,20 @@ VALID_COLUMNS = ['Filename', 'Corrected_PROGRAM', 'Corrected_NETWORK'] -# AWS Keys -AWS_ACCESS_KEY_ID="ASIAZTMXAXNXBHQKW32L" -AWS_SECRET_ACCESS_KEY="IozsUqPsyXu2klzNMU42CDZ1ZtGz8Yo29mhuJRaQ" -AWS_SESSION_TOKEN="IQoJb3JpZ2luX2VjEDsaCXVzLWVhc3QtMiJHMEUCIFkt/LwV1X3UmAyNi25xrDIL+mdaV7eV3U03VUQPc0o/AiEA4R54Gyd4d/55N5P0a84zzHiPUuejjyEMF6ARAk0i2fgqjAMI5P//////////ARAAGgw2NjAxMzEwNjg3ODIiDC+Kl0HOGnW+Vct/FyrgAg+t9ssFLXSj0VFXfVud6jLJnIAChoXSikj59wmU1hX3vcsN699t4UT5vfAjIwZ3emEPrfaKnuvcdC4ev2V4fT3IxRS8FLu6eUKqImF8ZkATIxg0nYpPqx+92CBvcT1gjLv4Y/1oorb0GxLcTiZLBY6IrCbrr3bDyexU+fgqoGlVFhaqfVPotvXy5+s0RbFWF7wQ0IanZ4eX5ZTFUXZMJh0LBYSTIuas6oDF2z/k/De48NKOWEZCmVxIG/ykyIbBujSICetKKX1qGtB74NrxSfhERGfof4DjZl5q4SUEouTRVb5aGk779lqdKjTXjmGyWyryzUu9jJfqyZiC+ymcIliK4aTn4W88/prqx3p1mq+/T1e5uMhtXi/9JURT3P/avOVz34j5ohggUxqvl2TBASKaUGLrOIBoJzoMevgodoWy1TleJe7xJWBOS5nq5knxMOreHjuImgNofmiRuSNbayowpqTzswY6pgFwPrVOCxqQyxXwNsCEB27me72Ez3splxi+JYiYyaxkANoftb6d1/JyDBprnsact7hvD3mPVUhjdvsh92o5+F1lR0usiyCqrvUrXgpYKdBh6jIJh2kcNDSYwoGMHI2kVb2G53GWe1Vb06cS79Ei3QVurxVe7SVYVqd+bpicuYDdcrvdzCeF2E+9P9qEQa1/4B6ei+0O+lao4WBjvy7ZyuAlwaUUgsw0" # File Paths LOCAL_PATH = 'data/test/' # Replace with local # S3 Settings S3_CLIENT = boto3.client('s3', - region_name="us-east-2", - aws_access_key_id=AWS_ACCESS_KEY_ID, - aws_secret_access_key=AWS_SECRET_ACCESS_KEY, - aws_session_token=AWS_SESSION_TOKEN + region_name="us-east-2" ) BUCKET = "doczy-dev-infra-textract" PREFIX = "batches/batch_1/contract-text-file/" # replace with s3 path # Bedrock Settings BEDROCK_RUNTIME = boto3.client(service_name="bedrock-runtime", - region_name="us-east-1", - aws_access_key_id=AWS_ACCESS_KEY_ID, - aws_secret_access_key=AWS_SECRET_ACCESS_KEY, - aws_session_token=AWS_SESSION_TOKEN + region_name="us-east-1" ) MODEL_ID_CLAUDE3_HAIKU = 'anthropic.claude-3-haiku-20240307-v1:0' MODEL_ID_CLAUDE3_SONNET = 'anthropic.claude-3-sonnet-20240229-v1:0' diff --git a/textract-pipeline/src/lambda/prompt-orchestrator/index.py b/textract-pipeline/src/lambda/prompt-orchestrator/index.py index e447663..9df42d7 100644 --- a/textract-pipeline/src/lambda/prompt-orchestrator/index.py +++ b/textract-pipeline/src/lambda/prompt-orchestrator/index.py @@ -12,6 +12,7 @@ import logging from urllib.parse import unquote_plus from botocore.exceptions import ClientError import traceback +import file_processing # Global configuration variables S3_REGION = "us-east-2" @@ -117,6 +118,12 @@ def save_to_sf(dag_name, **kwargs): return payload +def parse_field_groups(field_groups): + # Split the string by comma and strip any surrounding whitespace + groups = [group.strip() for group in field_groups.split(',')] + # Convert the list to a tuple + return tuple(groups) + def lambda_handler(event, context): try: logger.info("Processing SQS event.") @@ -136,17 +143,75 @@ def lambda_handler(event, context): # Fetch batch_id from object tags batch_id = get_s3_object_tags(bucket_name, object_key).get('BatchId') logger.info(f"Batch ID: {batch_id}") - - # Query field groups from the database - prompt_config_query = f"SELECT DISTINCT GROUP_ID FROM DOCUMENT_LOGS WHERE DOCUMENT_ID = '{document_id}'" - # field_group_df = read_from_db(prompt_config_query) - # field_groups = tuple(field_group_df['GROUP_ID'].unique()) - field_groups = ('A','C') - logger.info(f"Field groups: {field_groups}") - # Prepare and process the prompt - main(bucket_name, object_key, field_groups) + # Get file name from s3 bucket object key + try: + logger.info(f"OBJECT KEY:: {object_key}") + # Get file name from s3 bucket object key and remove file extension + file_name = object_key.split('/')[-1].split('.')[0] + cursor = snowflake_conn.cursor() + result = cursor.execute(f"select group_id from stg.document_logs where batch_id= '{batch_id}' and document_id='{file_name}' limit 1;") + logger.info(f"QUERY BEING EXECUTED: select group_id from stg.document_logs where batch_id= '{batch_id}' and file_name='{file_name}' limit 1;") + + + for rec in result: + logger.info(f"QUERY RESULT::: {rec[0]}") + field_groups = rec[0] + + # logger.info(f"QUERY RESULT::: {cursor.fetchone()} & {cursor.fetchone()[0]}&fetching all {cursor.fetchall()}") + # field_groups = cursor.fetchone()[0] + # logger.info(f"RESULT FIELD GROUP:: {field_groups}") + # Check if field group is a tuple, if not convert it to tuple + field_groups = parse_field_groups(field_groups) + logger.info(f"RESULT FIELD GROUP after tuple conversion:: {field_groups}") + + except Exception as e: + logger.error(f"Error fetching field groups: {e} stopping execution") + exit() + field_groups = ('A','C') + + + # field_groups = ('A','C') + # field_groups = 'B' + # logger.info(f"Field groups: {field_groups}") + if field_groups == ('B',): + data = s3_client.get_object(Bucket=bucket_name, Key=object_key) + contents = data['Body'].read() + context = contents.decode("utf-8") + # Get the file name from the s3 object prefix + filename = object_key.split('/')[-1] + process_b_fields(filename, context, bucket_name,batch_id) + # Prepare and process the prompt if field groups are A and C or Just A or just C + elif field_groups == ('A','C'): + main(bucket_name, object_key, field_groups) + # Process scenario where all field groups are processed + elif field_groups == ('A','B','C'): + main(bucket_name, object_key, ('A','C')) + data = s3_client.get_object(Bucket=bucket_name, Key=object_key) + contents = data['Body'].read() + context = contents.decode("utf-8") + # Get the file name from the s3 object prefix + filename = object_key.split('/')[-1] + process_b_fields(filename, context, bucket_name,batch_id) + elif field_groups == ('A',): + main(bucket_name, object_key, field_groups) + elif field_groups == ('C',): + main(bucket_name, object_key, field_groups) + elif field_groups == ('B','A') or field_groups == ('B','C') or field_groups == ('A','B') or field_groups == ('C','B'): + main(bucket_name, object_key, ('A','C')) + data = s3_client.get_object(Bucket=bucket_name, Key=object_key) + contents = data['Body'].read() + context = contents.decode("utf-8") + # Get the file name from the s3 object prefix + filename = object_key.split('/')[-1] + process_b_fields(filename, context, bucket_name,batch_id) + else: + logger.error(f"Unsupported field groups: {field_groups}") + return { + 'statusCode': 500, + 'body': f"Unsupported field groups: {field_groups}" + } return { 'statusCode': 200, @@ -159,6 +224,33 @@ def lambda_handler(event, context): 'body': f"{str(e)} \n {traceback.format_exc()}" } + +def process_b_fields(filename, context, bucket, batch_id): + # Create tuple from filename and context + item = (filename, context) + df = file_processing.process_file(item) + # Adding batch_id to the last column of df + df['Batch_id'] = batch_id + csv_buf = StringIO() + df.to_csv(csv_buf, header=True, index=False) + csv_buf.seek(0) + timestamp = datetime.now().strftime("%Y%m%d%H%M%S") + file_name = f"results_{batch_id}_B_{timestamp}.csv" + s3_client.put_object(Bucket=bucket, Body=csv_buf.getvalue(), Key=f"final_output/{batch_id}/{file_name}") + csv_buf = StringIO() + df.to_csv(csv_buf, header=True, index=False) + csv_buf.seek(0) + + # Saving data to snowflake raw data ingestion bucket & calling stored proc + s3_client.put_object(Bucket=snowflake_ingestion_bucket, Body=csv_buf.getvalue(), Key=f"doczy_pipeline_output/{file_name}") + # cur = snowflake_conn.cursor() + # load_data_query = f"CALL LOAD_DOCZY_PIPELINE_RAW_OUTPUT_B('{file_name}')" + # cur.execute(load_data_query) + # log df shape as output + logger.info(f"B fields Processed file: {filename} with shape: {df.shape}") + + + def get_filename_from_path(full_path): return os.path.basename(full_path) @@ -434,7 +526,7 @@ def main(bucket, object_key, field_groups): df.to_csv(csv_buf, header=True, index=False) csv_buf.seek(0) timestamp = datetime.now().strftime("%Y%m%d%H%M%S") - file_name = f"results_{batch_id}_{timestamp}.csv" + file_name = f"results_{batch_id}_AC_{timestamp}.csv" s3_client.put_object(Bucket=bucket, Body=csv_buf.getvalue(), Key=f"final_output/{batch_id}/{file_name}") csv_buf = StringIO() df.to_csv(csv_buf, header=True, index=False)