Removed sender/receiver folder

This commit is contained in:
Pankaj Katariya
2024-02-14 18:46:23 +05:30
parent a7717613f8
commit 5d90d22d07
2 changed files with 0 additions and 461 deletions
-228
View File
@@ -1,228 +0,0 @@
import json
import logging
import boto3
import os
from configparser import ConfigParser
from botocore.exceptions import ClientError
# Initialize logger
logger = logging.getLogger(__name__)
logger.setLevel(logging.INFO)
# Initialize S3 & Textract clients
s3_client = boto3.client('s3')
textract_client = boto3.client('textract')
# Function to load configuration from S3
def load_config_from_s3(bucket_name, file_key):
# Download the config file from S3
response = s3_client.get_object(Bucket=bucket_name, Key=file_key)
config_content = response['Body'].read().decode('utf-8')
# Parse the config file
config_parser = ConfigParser()
config_parser.read_string(config_content)
# Convert the configuration to a dictionary
config_dict = {}
for section in config_parser.sections():
config_dict[section] = {key.upper(): value for key, value in config_parser.items(section)}
return config_dict
# Function to construct JSON path from S3 object path
def generate_output_json_path(src_folder,dest_folder, s3_object):
# split string to remove source folder path
txt = s3_object.split(src_folder)
if len(txt) == 2:
txt = txt[1]
else:
txt = txt[0]
# change file extension - remove .pdf and add .json
if txt.lower().endswith(".pdf"):
return dest_folder+txt[0:-4]+".json"
# Function to get Textract document analysis
def get_textract_document_analysis(job_id, textract_client):
# Initialize an empty list to store blocks
all_blocks = []
next_token = None
response = {}
flag = True
try:
while True:
if flag:
# Call the Textract API to get document analysis
response = textract_client.get_document_analysis(
JobId=job_id
)
flag = False
else:
# Call the Textract API to get document analysis
response = textract_client.get_document_analysis(
JobId=job_id,
NextToken=next_token
)
job_status = response["JobStatus"]
logger.info("Job %s status is %s.", job_id, job_status)
# Merge the blocks from the current response
all_blocks.extend(response.get('Blocks', []))
# Check if there are more blocks to retrieve
next_token = response.get('NextToken')
if not next_token:
logger.info("No more Textract response to retrieve")
break
except ClientError:
logger.exception("Couldn't get data for job %s.", job_id)
raise
else:
# Remove unnecessary keys from the last response
last_response = response.copy()
last_response.pop('Blocks', None)
last_response.pop('ResponseMetadata', None)
logger.info("Removed 'ResponseMetadata' key")
# Merge with {'Blocks': all_blocks}
final_response = {'Blocks': all_blocks}
final_response.update(last_response)
logger.info("Final Textract response is contructed")
return final_response
# Function to save Textract response to S3
def upload_response_to_s3(response, bucket_name, object_key, s3_client):
# Convert the response to JSON
response_json = json.dumps(response)
try:
# Upload the JSON response to S3
s3_client.put_object(
Bucket=bucket_name,
Key=object_key,
Body=response_json,
ContentType='application/json'
)
logger.info(f"Saved analysis response to S3: {object_key}")
except Exception as e:
logger.error(f"Error saving response to S3 - {object_key}: {str(e)}")
# Function to move file within S3
def move_file_within_s3(source_bucket, source_path, destination_path):
try:
# Copy the file to the destination folder
s3_client.copy_object(Bucket=source_bucket, CopySource={'Bucket': source_bucket, 'Key': source_path}, Key=destination_path)
# Delete the file from the source folder
s3_client.delete_object(Bucket=source_bucket, Key=source_path)
logger.info(f"File moved from {source_path} to {destination_path}")
except Exception as e:
logger.error(f"Error moving file: {e}")
# AWS Lambda handler function
def lambda_handler(event, context):
logger.info('## ENVIRONMENT VARIABLES\r' + str(os.environ))
# Read environment variables
property_file_path = os.environ.get('PROPERTY_FILE_S3_PATH', '')
batch_id = os.environ.get('BATCH_ID', '')
# Validate environment variables
file_path_array = property_file_path.split("/")
if len(file_path_array) > 1:
# Extract BUCKET_NAME and config_file_path
S3_BUCKET_NAME = file_path_array[0]
CONFIG_FILE_PATH = "/".join(file_path_array[1:])
logger.info(f'S3_BUCKET_NAME: {S3_BUCKET_NAME}')
logger.info(f'CONFIG_FILE_PATH: {CONFIG_FILE_PATH}')
# Load config file
config_dict = load_config_from_s3(S3_BUCKET_NAME, CONFIG_FILE_PATH)
logger.info('## CONFIG DICTIONARY\r' + str(config_dict))
STAGING_LOCATION = config_dict['FOLDER_LOCATIONS']['STAGING_LOCATION'].format(batch_id)
OUTPUT_LOCATION = config_dict['FOLDER_LOCATIONS']['OUTPUT_LOCATION'].format(batch_id)
PROCESSED_LOCATION = config_dict['FOLDER_LOCATIONS']['PROCESSED_LOCATION'].format(batch_id)
UNPROCESSED_LOCATION = config_dict['FOLDER_LOCATIONS']['UNPROCESSED_LOCATION'].format(batch_id)
logger.info('STAGGING_LOCATION: ' + STAGING_LOCATION)
logger.info('OUTPUT_LOCATION: ' + OUTPUT_LOCATION)
logger.info('PROCESSED_LOCATION: ' + PROCESSED_LOCATION)
logger.info('UNPROCESSED_LOCATION: ' + UNPROCESSED_LOCATION)
# Process each message from the SQS event
for record in event['Records']:
# Extract the message body from the record
record_body = json.loads(record['body'])
message_body = json.loads(record_body['Message'])
logger.info('MESSAGE_BODY: ' + str(message_body))
try:
# Extract relevant information from the message body
job_id = message_body.get('JobId')
document_location = message_body.get('DocumentLocation')
s3_object_name = document_location.get('S3ObjectName')
# Check if the status is "SUCCEEDED"
if message_body.get('Status') == 'SUCCEEDED':
# Call the function to get document analysis using Textract
document_analysis = get_textract_document_analysis(job_id, textract_client)
# Save the document analysis response to S3
s3_object_key = generate_output_json_path(STAGING_LOCATION,OUTPUT_LOCATION, s3_object_name)
upload_response_to_s3(document_analysis, S3_BUCKET_NAME, s3_object_key, s3_client)
# Construct the destination paths
destination_path = PROCESSED_LOCATION + s3_object_name.replace(STAGING_LOCATION,"")
# Move file to processed folder
move_file_within_s3(S3_BUCKET_NAME, s3_object_name, destination_path)
success_message = 'Processed file '+ str(s3_object_key)
logger.info(success_message)
return {
'statusCode': 200,
'body': success_message
}
else:
error_message = f"Skipping message with JobId {message_body.get('JobId')} as Status is not 'SUCCEEDED'"
logger.info(error_message)
# Construct the destination paths
destination_path = UNPROCESSED_LOCATION + s3_object_name.replace(STAGING_LOCATION,"")
# Move file to unprocessed folder
move_file_within_s3(S3_BUCKET_NAME, s3_object_name, destination_path)
return {
'statusCode': 500,
'body': error_message
}
except Exception as e:
error_message = f"Error processing message with JobId {message_body.get('JobId')}: {str(e)}"
logger.error(error_message)
return {
'statusCode': 500,
'body': error_message
}
else:
error_message = 'Incorrect value for ENVIRONMENT VARIABLES: PROPERTY_FILE_S3_PATH\r' + str(property_file_path)
logger.error(error_message)
return {
'statusCode': 500,
'body': error_message
}
-233
View File
@@ -1,233 +0,0 @@
import boto3
import time
from configparser import ConfigParser
import logging
import os
from botocore.exceptions import ClientError
# Initialize logger
logger = logging.getLogger(__name__)
logger.setLevel(logging.INFO)
# Initialize S3 & Textract clients
s3_client = boto3.client('s3')
textract_client = boto3.client('textract')
# Function to generate a Unix timestamp
def generate_unix_timestamp():
# Get the current time in seconds since the epoch
unix_timestamp = int(time.time())
return unix_timestamp
# Function to retrieve configuration values from S3
def load_config_from_s3(bucket_name, file_key):
# Download the config file from S3
response = s3_client.get_object(Bucket=bucket_name, Key=file_key)
config_content = response['Body'].read().decode('utf-8')
# Parse the config file
config_parser = ConfigParser()
config_parser.read_string(config_content)
# Convert the configuration to a dictionary
config_dict = {}
for section in config_parser.sections():
config_dict[section] = {key.upper(): value for key, value in config_parser.items(section)}
return config_dict
# Function to move a file from source to destination in S3
def move_file_within_s3(source_bucket, source_key, destination_key):
try:
# Copy the file to the destination folder
s3_client.copy_object(Bucket=source_bucket, CopySource={'Bucket': source_bucket, 'Key': source_key}, Key=destination_key)
# Delete the file from the source folder
s3_client.delete_object(Bucket=source_bucket, Key=source_key)
logger.info(f"File moved from {source_key} to {destination_key}")
except ClientError as e:
logger.error(f"Error moving file: {e}")
except Exception as e:
logger.error(f"Error moving file: {e}")
# Function to get a list of PDF files in a given S3 folder
def get_pdf_files_list_from_s3(source_bucket, source_folder):
file_list = []
# List S3 Object & iterate (as per max files allowed)
s3_list_response = s3_client.list_objects_v2(Bucket=source_bucket, Prefix=source_folder)
if s3_list_response and s3_list_response['ResponseMetadata']['HTTPStatusCode'] == 200 and s3_list_response['KeyCount'] != 0:
objects = s3_list_response['Contents']
for s3_object in objects:
# Skip non-PDF files
if not s3_object['Key'].lower().endswith('.pdf'):
continue
file_list.append(s3_object['Key'])
return file_list
def start_textract_analysis_job(
bucket_name,
document_file_name,
analysis_feature_type,
sns_topic_arn,
sns_role_arn,
):
try:
# Define the parameters for the start_document_analysis API
start_document_analysis_params = {
'DocumentLocation': {
'S3Object': {
'Bucket': bucket_name,
'Name': document_file_name
}
},
'FeatureTypes': analysis_feature_type, # Customize based on requirements
'ClientRequestToken': 'unique-token-'+str(generate_unix_timestamp()), # Use a unique token for each request
'JobTag': 'healthcare-contract', # Use a tag to identify your job
'NotificationChannel': {
'SNSTopicArn': sns_topic_arn,
'RoleArn': sns_role_arn # Role to allow Textract service to notify SNS topic when response is ready
}
}
logger.info('start_document_analysis_params ' + str(start_document_analysis_params))
# Send the request to start document analysis
textract_response = textract_client.start_document_analysis(**start_document_analysis_params)
job_id = textract_response["JobId"]
logger.info(
"Started text analysis job %s on %s.", job_id, document_file_name
)
except ClientError:
logger.exception("Couldn't analyze text in %s.", document_file_name)
raise
else:
return job_id
# AWS Lambda handler function
def lambda_handler(event, context):
try:
# Extract AWS account ID and region from the Lambda ARN
aws_account_id = context.invoked_function_arn.split(":")[4]
aws_region = context.invoked_function_arn.split(":")[3]
logger.info('## ENVIRONMENT VARIABLES\r' + str(os.environ))
# Read environment variables
property_file_path = os.environ.get('PROPERTY_FILE_S3_PATH', '')
batch_id = os.environ.get('BATCH_ID', '')
# Read config.properties
file_path_array = property_file_path.split("/")
# Valid if file_path_array has more than 2 elements
if len(file_path_array) > 1:
# Extract BUCKET_NAME and config_file_path
S3_BUCKET_NAME = file_path_array[0]
CONFIG_FILE_PATH = "/".join(file_path_array[1:])
logger.info(f'S3_BUCKET_NAME: {S3_BUCKET_NAME}')
logger.info(f'CONFIG_FILE_PATH: {CONFIG_FILE_PATH}')
# Load config file
config_dict = load_config_from_s3(S3_BUCKET_NAME, CONFIG_FILE_PATH)
logger.info('## CONFIG DICTIONARY\r' + str(config_dict))
# Extract configuration values
SOURCE_LOCATION = config_dict['FOLDER_LOCATIONS']['SOURCE_LOCATION'].format(batch_id) # SOURCE_LOCATION
STAGING_LOCATION = config_dict['FOLDER_LOCATIONS']['STAGING_LOCATION'].format(batch_id)
ANALYSIS_FEATURE_TYPE = config_dict['OTHERS']['ANALYSIS_FEATURE_TYPE'].split(",") # Analysis FeatureType
SENDER_MAX_FILES = int(config_dict['OTHERS']['SENDER_MAX_FILES'])
SNS_TOPIC_ARN = config_dict['RESOURCES']['SNS_TOPIC_ARN'].replace("{aws_region}",aws_region).replace("{aws_account_id}",aws_account_id)
TEXTRACT_ROLE_ARN = config_dict['RESOURCES']['TEXTRACT_ROLE_ARN'].replace("{aws_account_id}",aws_account_id) # Textract IAM Role ARN to publish to SNS
logger.info('SOURCE_LOCATION: ' + SOURCE_LOCATION)
logger.info('STAGING_LOCATION: ' + STAGING_LOCATION)
logger.info('ANALYSIS_FEATURE_TYPE: ' + str(ANALYSIS_FEATURE_TYPE))
logger.info('SNS_TOPIC_ARN: ' + SNS_TOPIC_ARN)
logger.info('TEXTRACT_ROLE_ARN: ' + TEXTRACT_ROLE_ARN)
logger.info('SENDER_MAX_FILES: ' + str(SENDER_MAX_FILES))
# List S3 Object & iterate (as per max files allowed)
files_list = get_pdf_files_list_from_s3(S3_BUCKET_NAME,SOURCE_LOCATION)
if len(files_list):
# File count
file_count = 0
for s3_file_key in files_list:
# Construct the source and destination paths
source_path = s3_file_key
destination_path = STAGING_LOCATION + source_path.replace(SOURCE_LOCATION,"")
# Move file to stagging
move_file_within_s3(S3_BUCKET_NAME, source_path, destination_path)
# Start Textract analysis job
job_id = start_textract_analysis_job (
S3_BUCKET_NAME,
destination_path,
ANALYSIS_FEATURE_TYPE,
SNS_TOPIC_ARN,
TEXTRACT_ROLE_ARN,
)
file_count = file_count + 1
logger.info(str(file_count) + '. ' + str(s3_file_key) + " Job Id: " + str(job_id))
if SENDER_MAX_FILES == file_count:
break
success_message = 'Total files sent to textract : '+ str(file_count)
logger.info(success_message)
return {
'statusCode': 200,
'body': success_message
}
else:
message = 'No files found'
logger.error(message)
return {
'statusCode': 500,
'body': message
}
else:
error_message = 'Incorrect value for ENVIRONMENT VARIABLES: PROPERTY_FILE_S3_PATH\r' + str(property_file_path)
logger.error(error_message)
return {
'statusCode': 500,
'body': error_message
}
except ClientError as e:
# Handle specific Textract client errors
error_message = f"Error in Textract operation: {e}"
logger.error(error_message)
return {
'statusCode': 500,
'body': error_message
}
except Exception as e:
# Handle other exceptions
error_message = f"Unexpected error: {e}"
logger.error(error_message)
return {
'statusCode': 500,
'body': error_message
}