0% found this document useful (0 votes)
2 views29 pages

My code

The document contains Java code for configuring asynchronous task execution in a Spring application, including the setup of thread pools and a Kafka consumer for processing files. It defines a FileProcessor service that handles file processing tasks asynchronously, updating file states and managing exceptions. Additionally, it includes methods for success and failure callbacks based on the processing results, utilizing various services and repositories for data management.

Uploaded by

Anubhav Singh
Copyright
© All Rights Reserved
We take content rights seriously. If you suspect this is your content, claim it here.
Available Formats
Download as PDF, TXT or read online on Scribd
0% found this document useful (0 votes)
2 views29 pages

My code

The document contains Java code for configuring asynchronous task execution in a Spring application, including the setup of thread pools and a Kafka consumer for processing files. It defines a FileProcessor service that handles file processing tasks asynchronously, updating file states and managing exceptions. Additionally, it includes methods for success and failure callbacks based on the processing results, utilizing various services and repositories for data management.

Uploaded by

Anubhav Singh
Copyright
© All Rights Reserved
We take content rights seriously. If you suspect this is your content, claim it here.
Available Formats
Download as PDF, TXT or read online on Scribd

package [Link].

config;

import [Link];
import [Link].slf4j.Slf4j;
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];

@Configuration
@EnableAsync
@Slf4j
public class AsyncTaskExecutorConfig {

@Bean(name = "taskExecutor")
public Executor taskExecutor() {
final ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
[Link](8);
[Link](16);
[Link](200);
[Link]("ClassificationAsyncThread-");
[Link](true);
[Link](60);
[Link](true);
[Link](30);
[Link](new [Link]());
[Link]();
[Link]("Async Task Executor initialized successfully.");
return executor;

}
@Bean(name = "heavyPool")
public Executor taskExecutor() {
final ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
[Link](4);
[Link](4);
[Link](100);
[Link]("ClassificationAsyncThread-");
[Link](true);
[Link](60);
[Link](true);
[Link](30);
[Link](new [Link]());
[Link]();
[Link]("Async Task Executor initialized successfully.");
return executor;

@Bean
public ReentrantLock getReentrantLock() {
return new ReentrantLock(true);
}

package [Link];

import [Link];
import [Link];
import [Link];
import [Link];

import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];

import [Link];
import [Link];

import [Link].slf4j.Slf4j;

@Service
@Slf4j
@ConfigurationProperties(prefix = "[Link]")
public class FileConsumerVO {

private final FileProcessor fileProcessor;

private KafkaListenerEndpointRegistry registry;

// private Executor executor;

@Autowired
public FileConsumerVO(FileProcessor fileProcessor,
KafkaListenerEndpointRegistry registry) {
[Link] = fileProcessor;
[Link] = registry;
}
@KafkaListener(
topics = "${[Link]}",
groupId = "${[Link]}",
containerFactory = "kafkaListenerContainerFactoryvo"
)
public void consume(@Payload KafkaRequest kafkaRequest, Acknowledgment
acknowledgment) {
[Link](kafkaRequest);
[Link]();
[Link]("Acknowledgement sent for RequestId: {}",
[Link]().getRequestId());
}
}

package [Link];

import [Link];
import [Link];
import [Link];
import [Link];
import [Link];

import [Link];
import [Link].*;
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];

import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];

import [Link].slf4j.Slf4j;

@Service
@Slf4j
public class FileProcessor {

private final UploadService uploadService;


private final MongoClassificationRepo repo;
private final MasterCallBackUtil callBack;
private final SliceClassifyService sliceClassifyService;
private final FileProcessorStateRepository fileProcessorStateRepository;
private final CommonUtils commonUtils;
private final ValidationService validationService;
private ObjectMapper objectMapper;
private final RedisDataService redisDataService;
private String docValidationState = "";

@Autowired
public FileProcessor(UploadService uploadService, MasterCallBackUtil callBack,
MongoClassificationRepo repo,
SliceClassifyService sliceClassifyService, FileProcessorStateRepository
fileProcessorStateRepository,
CommonUtils commonUtils, ValidationService validationService,
RedisDataService redisDataService) {
[Link] = uploadService;
[Link] = callBack;
[Link] = repo;
[Link] = sliceClassifyService;
[Link] = fileProcessorStateRepository;
[Link] = commonUtils;
[Link] = validationService;
[Link] = redisDataService;
}

@Async("taskExecutor")
public void process(KafkaRequest kafkaRequest) {
[Link]("Listener thread async: {}", [Link]().getName());
Timestamp startTimestamp = [Link]([Link]());
[Link]("FileProcessor startTimestamp for RequestId: {} is {}",
[Link]().getRequestId(),
startTimestamp);
List<SliceAndClassifiedFileDetails> claimUploadFileDetailsList = new ArrayList<>();
List<Exception> exceptionList = new ArrayList<>();
IAssistAckResponse iAssistAckResponse = null;
boolean errorFlag = false;
Boolean updateflag = false;
String wholeDocValid = "";
try {
[Link]("Inside FileProcessor process method");
// Fetch Redis cache data
List<ClassificationMasterEntity> classificationMasterEntityList = redisDataService
.findAllClassificationMasterEntity();
List<FlagDetailsEntity> flagDetailsEntitiesList = [Link]();
[Link]("In FileProcessor : " + [Link](flagDetailsEntitiesList));
updateFileStateByFileId(kafkaRequest,
[Link].IN_PROGRESS.name(), "N");
[Link]("File state updated to INPROGRESS");
[Link]([Link](),
[Link](), uploadService, classificationMasterEntityList,
flagDetailsEntitiesList,
kafkaRequest.getS3uploadpath(), kafkaRequest.getS3BucketDetails(),
[Link]());

} catch (Exception exception) {


[Link]("Exception occurred during file processing : {}", [Link]());
[Link]("Stack Trace occurred during file processing : {}",
[Link](exception));
handleProcessingException(exception, claimUploadFileDetailsList, exceptionList);
errorFlag = true;
} finally {
[Link]("Finally FileProcessor process method");
[Link]("exceptionList: {}, size: {}", exceptionList, [Link]());
callSuccessFailureCallBack(kafkaRequest, wholeDocValid, exceptionList, updateflag);
if (null != claimUploadFileDetailsList) {
[Link]();
claimUploadFileDetailsList = null;
}
if (![Link]()) {
docValidationState = null;
}
}
}

private void updateFileProcessSkipRetry(KafkaRequest kafkaRequest, String


fileProcessedStatus, String filestatus,
String active, String wholeDocValid) {
[Link]("Started updateFileProcessedByFileId");
[Link]("wholeDocValid: {}", wholeDocValid);
Optional<List<FileProcessorStateEntity>> fileProcessorStateList =
fileProcessorStateRepository
.findByClientIdAndUsecaseIdAndUploadIdAndRequestIdAndFileId([Link]-
questDto().getClientId(),
[Link]().getUseCaseId(),
[Link]().getInputData().getUploadId(),
[Link]().getRequestId(), [Link]().getFileId());

[Link](states -> {
if (![Link]()) {
FileProcessorStateEntity stateEntity = [Link](0);
[Link]([Link]());
[Link]("KafkaFileConsumer");
[Link](fileProcessedStatus);
[Link](filestatus);
[Link](active);
if (null != wholeDocValid && ![Link]()
&& ValidationConstants.NOT_VALID_DOC.equalsIgnoreCase(wholeDocValid)) {
[Link](SlicingClassificationConsumerConstants.RETRY_COUNT);
[Link]([Link]);
[Link]([Link]);
}
[Link](stateEntity);
}
});
[Link]("Ended updateFileProcessedByFileId");
}

public void updateFileStateByFileId(KafkaRequest kafkaRequest, String status, String


fileProcessedStatus) {
[Link]("Started updateFileStateByFileId");
Optional<List<FileProcessorStateEntity>> fileProcessorStateList =
fileProcessorStateRepository
.findByClientIdAndUsecaseIdAndUploadIdAndRequestIdAndFileId([Link]-
questDto().getClientId(),
[Link]().getUseCaseId(),
[Link]().getInputData().getUploadId(),
[Link]().getRequestId(), [Link]().getFileId());
[Link](states -> {
if (![Link]()) {
FileProcessorStateEntity stateEntity = [Link](0);
[Link](status);
if
([Link]([Link](
))) {
[Link]([Link]);
}
[Link]([Link]());
[Link]("KafkaFileConsumer");
[Link](fileProcessedStatus);
[Link](stateEntity);
}
});
[Link]("Ended updateFileStateByFileId");
}

private void handleProcessingException(Exception e, List<SliceAndClassifiedFileDetails>


claimUploadFileDetailsList,
List<Exception> exceptionList) {
Throwable rootCause = [Link](e);
if (rootCause instanceof ClassificationCustomException) {
claimUploadFileDetailsList
.add(((ClassificationCustomException) rootCause).getSliceAndClassifiedFileDetails());
[Link](e);
} else if (rootCause instanceof ValidationCustomException) {
docValidationState = ((ValidationCustomException)
rootCause).getSliceAndClassifiedFileDetails()
.getWholeDocValidation();
[Link](
((ValidationCustomException)
[Link](e)).getSliceAndClassifiedFileDetails());
[Link](e);
} else {
SliceAndClassifiedFileDetails errorDetails = new SliceAndClassifiedFileDetails();

[Link](ErrorCodes.CLASSIFICATION_INTERNAL_SERVER_ERROR.getError-
Code());

[Link](ErrorCodes.CLASSIFICATION_INTERNAL_SERVER_ERROR.getEr-
rorMessage());
[Link](errorDetails);
[Link](e);
}
}

private void callSuccessFailureCallBack(KafkaRequest kafkaRequest, String wholeDocValid,


List<Exception> exceptionList, Boolean updateflag) {
[Link]("Started callSuccessFailureCallBack");
[Link]("callSuccessFailureCallBack exceptionList: {}, size: {}", exceptionList,
[Link]());
Optional<List<FileProcessorStateEntity>> fileProcessedList = fileProcessorStateRepository
.findByClientIdAndUsecaseIdAndUploadIdAndRequestId([Link]().g
etClientId(),
[Link]().getUseCaseId(),
[Link]().getInputData().getUploadId(),
[Link]().getRequestId());
if ([Link]() && ![Link]().isEmpty() &&
[Link]().size() > 0) {
int totalFileCount = [Link]().size();
[Link]("totalFileCount {} : ", totalFileCount);
int processedFileCount = (int) [Link]().stream()
.filter(s ->
[Link]([Link]())).count();
[Link]("processedFileCount {} : ", processedFileCount);
if (totalFileCount == processedFileCount) {
int failedFileCount = (int) [Link]().stream()
.filter(s -> [Link]()
.equalsIgnoreCase([Link]()))
.count();
[Link]("failedFileCount {} : ", failedFileCount);
if (failedFileCount > 0) {
handleFailureCallbackBasedOnValidation(kafkaRequest, wholeDocValid, exceptionList,
updateflag);
} else {
handleSuccessCallbackBasedOnValidation(kafkaRequest, wholeDocValid, exceptionList);
}
}
}
[Link]("Ended callSuccessFailureCallBack");
}

private void handleSuccessCallbackBasedOnValidation(KafkaRequest kafkaRequest, String


wholeDocValid,
List<Exception> exceptionList) {
[Link]("Started handleCallbackBasedOnValidation ");
[Link]("Started general callback ");
BaseSyncAsyncResponse response =
createSuccessCallbackResponse([Link]());
[Link](response, [Link]());
[Link]("Ended handleCallbackBasedOnValidation ");
}

private void handleFailureCallbackBasedOnValidation(KafkaRequest kafkaRequest, String


wholeDocValid,
List<Exception> exceptionList, Boolean updateflag) {
[Link]("Started handleFailureCallbackBasedOnValidation ");
[Link]("inside handleFailureCallbackBasedOnValidation: exceptionList: {}, size: {}",
exceptionList, [Link]());
BaseSyncAsyncResponse baseSyncAsyncResponse =
createErrorCallback([Link](),
ErrorCodes.CLASSIFICATION_INTERNAL_SERVER_ERROR);
[Link](baseSyncAsyncResponse,
[Link](baseSyncAsyncResponse,
[Link]());
[Link]("Ended handleFailureCallbackBasedOnValidation ");
}

public BaseSyncAsyncResponse createErrorCallback(SliceAndClassifyRequestDto requestDto,


ErrorCodes errorCodes) {
[Link]("Started createErrorCallback");
BaseSyncAsyncResponse baseSyncAsyncResponse = new BaseSyncAsyncResponse();
BaseErrorResponse baseErrorResponse = new BaseErrorResponse();
[Link]([Link]());
[Link]([Link]());
if ([Link]() != null) {
[Link]([Link]());
} else {
[Link](null);
}
[Link]([Link]().getUploadId());
[Link]([Link]());
[Link]([Link]());
[Link]([Link]());
[Link](baseErrorResponse);
return baseSyncAsyncResponse;
}

private static BaseSyncAsyncResponse createDCSuccesscallback(SliceAndClassifyRequestDto


requestDto) {
BaseSyncAsyncResponse baseSyncAsyncResponse = new BaseSyncAsyncResponse();
BaseSuccessResponse baseSuccessResponse = new BaseSuccessResponse();
[Link]([Link]().getUploadId());
[Link]([Link]());
[Link]([Link]());
[Link]([Link]());
[Link]([Link]().getCategory());
[Link]([Link]().getSubcategory());
[Link]([Link]().getFlowname());
baseSuccessResponse
.setResponseMessage(SlicingClassificationConsumerConstants.SUCCESS_RESPON-
SE_MESSAGE_SLICING_CLASSIFICATION);

[Link]([Link]-
CESS_STATUS_CODE);
baseSuccessResponse
.setDisplayMessage(SlicingClassificationConsumerConstants.SUCCESS_DISPLAY_MES-
SAGE_SLICING_CLASSIFICATION);
if (null != [Link]().getTags() && !
[Link]().getTags().isEmpty()) {
[Link]([Link]().getTags());
}
if (null != [Link]().getMetadata()) {
[Link](createMetaData(requestDto));
}
[Link]([Link]().getAsync());
[Link]([Link]().getPolicyNo());
[Link]([Link]().getHealthCardId());
[Link]([Link]().getSource());
[Link](baseSuccessResponse);
return baseSyncAsyncResponse;
}

private BaseSyncAsyncResponse createSuccessCallbackResponse(SliceAndClassifyRequestDto


requestDto) {
BaseSyncAsyncResponse response = new BaseSyncAsyncResponse();
BaseSuccessResponse successResponse = new BaseSuccessResponse();
[Link]([Link]().getUploadId());
[Link]([Link]().getFlowname());
[Link]([Link]());
[Link]([Link]());
[Link]([Link]());
[Link]([Link]().getCategory());
[Link]([Link]().getSubcategory());
successResponse
.setResponseMessage(SlicingClassificationConsumerConstants.IASSIST_ACK_MES-
SAGE_SLICING_CLASSIFICATION);

[Link](SlicingClassificationConsumerConstants.ACCEPTED_STA-
TUS_CODE);
successResponse
.setDisplayMessage(SlicingClassificationConsumerConstants.IASSIST_ACK_MES-
SAGE_SLICING_CLASSIFICATION);
if (null != [Link]().getTags() && !
[Link]().getTags().isEmpty()) {
[Link]([Link]().getTags());
}
[Link](createMetaData(requestDto));
[Link]([Link]().getAsync());
[Link]([Link]().getPolicyNo());
[Link]([Link]().getHealthCardId());
[Link]([Link]().getSource());
[Link](successResponse);
return response;
}

private static BaseSyncAsyncResponse createSuccesscallback(ValidationRequestDto


vrequestDto) {
BaseSyncAsyncResponse baseSyncAsyncResponse = new BaseSyncAsyncResponse();
BaseSuccessResponse baseSuccessResponse = new BaseSuccessResponse();
[Link]([Link]().getUploadId());
[Link]([Link]().getFlowname());
[Link]([Link]());
[Link]([Link]());
[Link]([Link]());
[Link]([Link]().getCategory());
[Link]([Link]().getSubcategory());

[Link]([Link]-
CESS_RESPONSE_MESSAGE_VALIDATION);

[Link]([Link]-
CESS_STATUS_CODE);

[Link]([Link]-
CESS_DISPLAY_MESSAGE_VALIDATION);
if (null != [Link]()) {
[Link]([Link]());
}
if (null != [Link]().getTags() && !
[Link]().getTags().isEmpty()) {
[Link]([Link]().getTags());
}
if (null != [Link]().getMetadata()) {
[Link](createMetaData(vrequestDto));
}
[Link]([Link]().getAsync());
[Link]([Link]().getPolicyNo());
[Link]([Link]().getHealthCardId());
[Link]([Link]().getSource());
[Link](baseSuccessResponse);

return baseSyncAsyncResponse;
}

private static MetaData createMetaData(ValidationRequestDto vrequestDto) {


MetaData metaData = new MetaData();
if (null != vrequestDto && null != [Link]()
&& null != [Link]().getMetadata()) {
if ([Link]().getMetadata().getProposalNumber() != null) {

[Link]([Link]().getMetadata().getProposalNumber());
}
if ([Link]().getMetadata().getCustomerId() != null) {
[Link]([Link]().getMetadata().getCustomerId());
}
if ([Link]().getMetadata().getClaimId() != null) {
[Link]([Link]().getMetadata().getClaimId());
}
if ([Link]().getMetadata().getPolicyNumber() != null) {
[Link]([Link]().getMetadata().getPolicyNumber());
}
}
return metaData;
}

private static MetaData createMetaData(SliceAndClassifyRequestDto requestDto) {


MetaData metaData = new MetaData();
if (null != requestDto && null != [Link]()
&& null != [Link]().getMetadata()) {
if ([Link]().getMetadata().getProposalNumber() != null) {

[Link]([Link]().getMetadata().getProposalNumber());
}
if ([Link]().getMetadata().getCustomerId() != null) {
[Link]([Link]().getMetadata().getCustomerId());
}
if ([Link]().getMetadata().getClaimId() != null) {
[Link]([Link]().getMetadata().getClaimId());
}
if ([Link]().getMetadata().getPolicyNumber() != null) {
[Link]([Link]().getMetadata().getPolicyNumber());
}
}
return metaData;
}

private List<ClassificationMasterEntity> redisFindAll() {


return [Link]();
}

package [Link];

import [Link];

import [Link].*;
import [Link];
import [Link];
import [Link];

import [Link];
import [Link];
import [Link];

import [Link].slf4j.Slf4j;

@Service
@Slf4j
public class SliceClassifyServiceImpl implements SliceClassifyService {

@Override
public void sliceAndClassify(FileModel file, SliceAndClassifyRequestDto requestDto,
UploadService uploadService,
List<ClassificationMasterEntity> classificationMasterEntityList,
List<FlagDetailsEntity>
flagDetailsEntityList, String s3uploadpath, S3BucketDetails s3BucketDetails, String
todayYearMonthDate)
throws Exception, ClassificationCustomException {

UploadDocUrl uploadDocUrl = new UploadDocUrl();


[Link]([Link]());
[Link]("Processing file with ID: {}", [Link]());

[Link]([Link]());

try {
[Link](uploadDocUrl, requestDto, [Link](),
classificationMasterEntityList,flagDetailsEntityList, s3uploadpath,
s3BucketDetails, todayYearMonthDate);
} catch (ClassificationCustomException exception) {
[Link]("Exception occurred during file slicing and classification: {}",
[Link](), exception);
[Link]("Stack Trace occurred during file processing : {}",
[Link](exception));
throw exception;
} catch (Exception exception) {
[Link]("Exception occurred during SliceClassifyServiceImpl: {}",
[Link](), exception);
[Link]("Stack Trace occurred during SliceClassifyServiceImpl : {}",
[Link](exception));
throw exception;
}

}
}

package [Link];

import [Link];
import [Link].*;
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];

@Service
public class UploadServiceImpl implements UploadService {

private AsyncUploadService asyncUploadService;

@Autowired
public UploadServiceImpl(AsyncUploadService asyncUploadService) {
[Link] = asyncUploadService;
}

@Override
public void sliceAndClassify(UploadDocUrl uploadDocUrl,
SliceAndClassifyRequestDto requestDto,
String fileId,
List<ClassificationMasterEntity>
classificationMasterEntityList,

List<FlagDetailsEntity>flagDetailsEntityList, String s3uploadpath,


S3BucketDetails s3BucketDetails, String
todayYearMonthDate) throws Exception {
[Link](uploadDocUrl, requestDto, fileId,
classificationMasterEntityList,flagDetailsEntityList,
s3uploadpath, s3BucketDetails, todayYearMonthDate);
}
}

package [Link];

import [Link];
import [Link];
import [Link];
import [Link];
import [Link].*;

import [Link];
import [Link].*;
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link].*;
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];

import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];

import [Link].slf4j.Slf4j;
import [Link];
import [Link].S3Exception;
@Slf4j
@Service
public class AsyncUploadService {

private AWSFileUploadService awsFileUploadService;


private FileServices fileServices;
private IAssistService iAssistService;
private ExceptionUtil exceptionUtil;
private ImageStructuringService imageStructuringService;
private PreLoadProperty preLoadProperty;
private SliceAndConvertToImageService fileSliceAndConvertToImageImpl;
private MongoClassificationRepo repo;
private FileProcessorStateRepository fileProcessorStateRepository;
private CommonUtils commonUtils;

@Autowired
public AsyncUploadService(AWSFileUploadService awsFileUploadService, FileServices
fileServices,
IAssistService iAssistService, ExceptionUtil exceptionUtil,
ImageStructuringService imageStructuringService,
PreLoadProperty preLoadProperty,
@Qualifier("fileSliceAndConvertToImageImpl") SliceAndConvertToImageService
fileSliceAndConvertToImageImpl,
MongoClassificationRepo repo, FileProcessorStateRepository
fileProcessorStateRepository, CommonUtils commonUtils) {
[Link] = awsFileUploadService;
[Link] = fileServices;
[Link] = iAssistService;
[Link] = exceptionUtil;
[Link] = imageStructuringService;
[Link] = preLoadProperty;
[Link] = fileSliceAndConvertToImageImpl;
[Link] = repo;
[Link] = fileProcessorStateRepository;
[Link] = commonUtils;
}

public void splitDocument(UploadDocUrl uploadDocUrl, SliceAndClassifyRequestDto


requestDto,
String fileId, List<ClassificationMasterEntity> classificationMasterEntityList,
List<FlagDetailsEntity> flagDetailsEntityList,
String s3uploadpath, S3BucketDetails s3BucketDetails, String
todayYearMonthDate)
throws Exception, ClassificationCustomException {
[Link]("Started splitDocument");
IAssistAckResponse iAssistAckResponse = null;
SliceAndClassifiedFileDetails sliceAndClassifiedFileDetails = new
SliceAndClassifiedFileDetails();
String originalFileUrl = [Link]();
String uriPath = [Link](originalFileUrl);
String partNo = fileId;

[Link]("uriPath: {}", uriPath);

Map<String, String> pm = [Link](partNo, originalFileUrl);


String fileNameWithExt = [Link](uriPath);
List<String> seprateFileAndExt = [Link]([Link]("\\."));
Map<String, String> fileType = [Link](partNo, fileNameWithExt);

[Link](pm);
[Link](fileType);
[Link](fileId);

convertToImage(partNo, originalFileUrl, seprateFileAndExt,


requestDto, fileId, sliceAndClassifiedFileDetails, classificationMasterEntityList,
s3BucketDetails, s3uploadpath, todayYearMonthDate);

[Link]("Ended splitDocument");
}

public void convertToImage(String partNo, String originalFileUrl,


List<String> seprateFileAndExt, SliceAndClassifyRequestDto requestDto, String
fileId,
SliceAndClassifiedFileDetails sliceAndClassifiedFileDetails,
List<ClassificationMasterEntity> classificationMasterEntityList, S3BucketDetails
s3BucketDetails,
String s3uploadpath, String todayYearMonthDate) throws Exception,
ClassificationCustomException {

[Link]("Inside convertToImage, PDF page slicing and conversion to be started");


IAssistAckResponse iAssistAckResponse = null;
List<ImageDetails> imageDetailsList = new ArrayList<>();
Map<List<ImageDetails>, Boolean> updatedResp = null;
int loopSize = 0;
try (InputStream inputStream = awsFileUploadService.downloadFileFromS3(originalFileUrl,
s3BucketDetails)) {
if
(SlicingClassificationConsumerConstants.FILE_TYPE_PDF.equals([Link](1))) {
processPdfFile(partNo, inputStream, seprateFileAndExt, imageDetailsList, s3BucketDetails,
s3uploadpath, todayYearMonthDate, loopSize, requestDto);
} else {
loopSize++;
// Process non-PDF files if needed
processNonPdfFile(partNo, inputStream, seprateFileAndExt, imageDetailsList,
s3BucketDetails, s3uploadpath, todayYearMonthDate, loopSize, requestDto);
}

iAssistAckResponse = [Link](imageDetailsList, originalFileUrl,


requestDto, fileId, classificationMasterEntityList, sliceAndClassifiedFileDetails);
[Link]("Converted the pdf into multiple jpeg");
[Link]("ACK_RECEIVED");
updateMongoDBForSuccess(sliceAndClassifiedFileDetails, requestDto, fileId,
iAssistAckResponse);
} catch (Exception e) {
[Link](requestDto,
[Link],
[Link](),
[Link], fileId);;
handleException(sliceAndClassifiedFileDetails, e, requestDto);
}
[Link]("Inside convertToImage, PDF page slicing and conversion to be Ended");
}

private void processPdfFile(String partNo, InputStream inputStream, List<String>


seprateFileAndExt,
List<ImageDetails> imageDetailsList, S3BucketDetails s3BucketDetails,
String s3uploadpath, String todayYearMonthDate, int loopSize,
SliceAndClassifyRequestDto requestDto) throws Exception {
[Link]("Splitting and conversion of PDF started with memory-optimized streaming");

// Use streaming PDF processor to avoid loading entire PDF into memory
try (StreamingPDFProcessor pdfProcessor = new StreamingPDFProcessor(inputStream)) {

Integer noOfPages = (loopSize == 0) ? [Link]() : 1;


PDFRenderer pdfRenderer = [Link]();

[Link]("PDF loaded with {} pages using streaming processor", noOfPages);

// Process pages one at a time to minimize memory usage


for (int page = 0; page < noOfPages; page++) {
PDPage pdPage = [Link](page);

ImageDetails imageDetails =
[Link](
partNo, page, inputStream, pdfRenderer, fileServices,
awsFileUploadService, seprateFileAndExt, exceptionUtil, s3BucketDetails,
s3uploadpath, pdPage, todayYearMonthDate, preLoadProperty, requestDto);
[Link]("Processed page {} - URL: {}", page, imageDetails != null ?
[Link]() : "null");

if (imageDetails != null) {
[Link](imageDetails);
}

} catch (IOException e) {
[Link]("IO Exception occurred during streaming PDF processing", e);
throw e;
} catch (Exception e) {
[Link]("Exception occurred during streaming PDF processing", e);
throw e;
}
}

private void processNonPdfFile(String partNo, InputStream inputStream, List<String>


seprateFileAndExt,
List<ImageDetails> imageDetailsList, S3BucketDetails s3BucketDetails,
String s3uploadpath, String todayYearMonthDate, int loopSize,
SliceAndClassifyRequestDto requestDto) throws Exception {
[Link]("Splitting and conversion of PDF started with memory-optimized streaming");

// Use streaming PDF processor to avoid loading entire PDF into memory
try {

Integer noOfPages = loopSize;

[Link]("PDF loaded with {} pages using streaming processor", noOfPages);

// Process pages one at a time to minimize memory usage


for (int page = 0; page < noOfPages; page++) {
ImageDetails imageDetails =
[Link](
partNo, page, inputStream, null, fileServices,
awsFileUploadService, seprateFileAndExt, exceptionUtil, s3BucketDetails,
s3uploadpath, null, todayYearMonthDate, preLoadProperty, requestDto);

[Link]("Processed page {} - URL: {}", page, imageDetails != null ?


[Link]() : "null");

if (imageDetails != null) {
[Link](imageDetails);
}
}

} catch (IOException e) {
[Link]("IO Exception occurred during streaming PDF processing", e);
throw e;
} catch (Exception e) {
[Link]("Exception occurred during streaming PDF processing", e);
throw e;
}
}

private byte[] convertToByte(InputStream input) {


try (ByteArrayOutputStream byteStream = new ByteArrayOutputStream()) {
int ch;
while ((ch = [Link]()) != -1) {
[Link](ch);
}
return [Link]();
} catch (IOException e) {
[Link]("Exception during byte conversion: {}", [Link](), e);
return new byte[0];
}
}
private void handleException(SliceAndClassifiedFileDetails sliceAndClassifiedFileDetails,
Exception e, SliceAndClassifyRequestDto requestDto) throws ClassificationCustomException {
[Link]("Exception occurred: {}", [Link](), e);

String errorCode = ErrorCodes.DEFAULT_ERROR.getErrorCode();


String errorMessage = ErrorCodes.DEFAULT_ERROR.getErrorMessage();

if (e instanceof S3Exception) {
errorCode = ErrorCodes.S3_EXCEPTION.getErrorCode();
errorMessage = ErrorCodes.S3_EXCEPTION.getErrorMessage();
} else if (e instanceof AbortedException) {
errorCode = ErrorCodes.ABORTEDEXCEPTION_ERROR.getErrorCode();
errorMessage = ErrorCodes.ABORTEDEXCEPTION_ERROR.getErrorMessage();
} else if (e instanceof IAssistBadGateway) {
errorCode = ErrorCodes.IASSIST_BAD_GATEWAY_ERROR.getErrorCode();
errorMessage = ErrorCodes.IASSIST_BAD_GATEWAY_ERROR.getErrorMessage();
} else if (e instanceof IAssistServiceUnavailable) {
errorCode = ErrorCodes.IASSIST_SERVICE_UNAVAILABLE.getErrorCode();
errorMessage = ErrorCodes.IASSIST_SERVICE_UNAVAILABLE.getErrorMessage();
} else if (e instanceof InternalServerError) {
errorCode = ErrorCodes.IASSIST_INTERNAL_SERVER_ERROR.getErrorCode();
errorMessage = ErrorCodes.IASSIST_INTERNAL_SERVER_ERROR.getErrorMessage();
} else if (e instanceof ImageConvException) {
errorCode = ErrorCodes.IMAGE_CONVERSION_ERROR.getErrorCode();
errorMessage = ((ImageConvException) e).getMessage() != null ? ((ImageConvException)
e).getMessage() : ErrorCodes.IMAGE_CONVERSION_ERROR.getErrorMessage();
} else if (e instanceof NoSuchElementException) {
errorCode = ErrorCodes.CLASSIFICATION_NOT_FOUND_DB.getErrorCode();
errorMessage = ErrorCodes.CLASSIFICATION_NOT_FOUND_DB.getErrorMessage();
}

[Link](errorCode);
[Link](errorMessage);
updateMongoDBForFailure(sliceAndClassifiedFileDetails, requestDto);
throw new ClassificationCustomException(sliceAndClassifiedFileDetails);
}

private void updateMongoDBForSuccess(SliceAndClassifiedFileDetails


sliceAndClassifiedFileDetails, SliceAndClassifyRequestDto requestDto, String fileId,
IAssistAckResponse iAssistAckResponse) {
List<SliceAndClassifiedFileDetails> claimUploadFileDetailsList = new ArrayList<>();

if (sliceAndClassifiedFileDetails != null) {
[Link](sliceAndClassifiedFileDetails);
}

if (![Link]()) {
ClassificationEntity classificationEntity = new ClassificationEntity();
[Link](claimUploadFileDetailsList);
[Link](requestDto, [Link](),
[Link]());
}

[Link](requestDto,
[Link],
[Link].WAITING_FOR_IASSIT_CALLBACK.name(),
[Link], sliceAndClassifiedFileDetails, fileId,
iAssistAckResponse);
}

private void updateMongoDBForFailure(SliceAndClassifiedFileDetails


sliceAndClassifiedFileDetails, SliceAndClassifyRequestDto requestDto) {
[Link]("Inside updateMongoDBForFailure");
List<SliceAndClassifiedFileDetails> claimUploadFileDetailsList = new ArrayList<>();

if (sliceAndClassifiedFileDetails != null) {
[Link](sliceAndClassifiedFileDetails);
}

if (![Link]()) {
ClassificationEntity classificationEntity = new ClassificationEntity();
[Link](claimUploadFileDetailsList);
[Link](requestDto, [Link](),
[Link]());
}
}
}

You might also like