Merged in feature/b_field_final (pull request #176)

Feature/b field final
This commit is contained in:
Umang Mistry
2024-06-28 22:47:46 +00:00
25 changed files with 1050 additions and 22 deletions
+29
View File
@@ -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;'
;
@@ -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)
);
@@ -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
);
@@ -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()
);
@@ -0,0 +1,6 @@
CREATE TABLE IF NOT EXISTS STG.REQUEST_SUBMISSION (
CLIENT_NAME VARCHAR,
BATCH_ID VARCHAR,
REQUEST_USERNAME VARCHAR,
REQUEST_DATETIME DATETIME
);
@@ -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;
$$;
@@ -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;
$$;
@@ -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)
);
@@ -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;
$$;
@@ -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
);
@@ -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;
$$;
@@ -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)
);
@@ -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;
@@ -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;
$$;
@@ -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;
$$;
@@ -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;
';
@@ -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;
$$;
@@ -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;
$$;
@@ -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
);
@@ -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;
$$;
@@ -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,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;
$$;
@@ -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
);
@@ -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'
@@ -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)