From 99b5e307789c3301672ed0111b42049051684155 Mon Sep 17 00:00:00 2001 From: unknown Date: Thu, 7 Feb 2019 21:20:05 +0530 Subject: [PATCH] Implementing parallelizing and paging of update loan summary --- .../LoanApplicationWriteServiceImpl.java | 124 ++++++++++++++++++ .../service/LoanReadPlatformService.java | 2 + .../service/LoanReadPlatformServiceImpl.java | 5 + .../service/ScheduledJobRunnerService.java | 2 + .../ScheduledJobRunnerServiceImpl.java | 89 ++++++++++++- 5 files changed, 217 insertions(+), 5 deletions(-) create mode 100644 fineract-provider/src/main/java/org/apache/fineract/portfolio/loanaccount/service/LoanApplicationWriteServiceImpl.java diff --git a/fineract-provider/src/main/java/org/apache/fineract/portfolio/loanaccount/service/LoanApplicationWriteServiceImpl.java b/fineract-provider/src/main/java/org/apache/fineract/portfolio/loanaccount/service/LoanApplicationWriteServiceImpl.java new file mode 100644 index 00000000000..9d859771350 --- /dev/null +++ b/fineract-provider/src/main/java/org/apache/fineract/portfolio/loanaccount/service/LoanApplicationWriteServiceImpl.java @@ -0,0 +1,124 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.fineract.portfolio.loanaccount.service; + +import org.apache.fineract.infrastructure.core.service.RoutingDataSourceServiceFactory; +import org.apache.fineract.infrastructure.core.service.ThreadLocalContextUtil; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.stereotype.Service; + +@Service +public class LoanApplicationWriteServiceImpl implements LoanApplicationWriteService { + + private final static Logger logger = LoggerFactory.getLogger(LoanApplicationWriteServiceImpl.class); + + private final RoutingDataSourceServiceFactory dataSourceServiceFactory; + + @Autowired + public LoanApplicationWriteServiceImpl(final RoutingDataSourceServiceFactory dataSourceServiceFactory){ + this.dataSourceServiceFactory = dataSourceServiceFactory; + } + + @Override + public void updateLoanSummaryDetails(Integer limit,Integer lastLoanId,String officeHierachy) { + final JdbcTemplate jdbcTemplate = new JdbcTemplate(this.dataSourceServiceFactory.determineDataSourceService().retrieveDataSource()); + + final StringBuilder updateSqlBuilder = new StringBuilder(900); + updateSqlBuilder.append("update m_loan "); + updateSqlBuilder.append("join ("); + updateSqlBuilder.append("SELECT ml.id AS loanId,"); + updateSqlBuilder.append("SUM(mr.principal_amount) as principal_disbursed_derived, "); + updateSqlBuilder.append("SUM(IFNULL(mr.principal_completed_derived,0)) as principal_repaid_derived, "); + updateSqlBuilder.append("SUM(IFNULL(mr.principal_writtenoff_derived,0)) as principal_writtenoff_derived,"); + updateSqlBuilder.append("SUM(IFNULL(mr.interest_amount,0)) as interest_charged_derived,"); + updateSqlBuilder.append("SUM(IFNULL(mr.interest_completed_derived,0)) as interest_repaid_derived,"); + updateSqlBuilder.append("SUM(IFNULL(mr.interest_waived_derived,0)) as interest_waived_derived,"); + updateSqlBuilder.append("SUM(IFNULL(mr.interest_writtenoff_derived,0)) as interest_writtenoff_derived,"); + updateSqlBuilder + .append("SUM(IFNULL(mr.fee_charges_amount,0)) + IFNULL((select SUM(lc.amount) from m_loan_charge lc where lc.loan_id=ml.id and lc.is_active=1 and lc.charge_time_enum=1),0) as fee_charges_charged_derived,"); + updateSqlBuilder + .append("SUM(IFNULL(mr.fee_charges_completed_derived,0)) + IFNULL((select SUM(lc.amount_paid_derived) from m_loan_charge lc where lc.loan_id=ml.id and lc.is_active=1 and lc.charge_time_enum=1),0) as fee_charges_repaid_derived,"); + updateSqlBuilder.append("SUM(IFNULL(mr.fee_charges_waived_derived,0)) as fee_charges_waived_derived,"); + updateSqlBuilder.append("SUM(IFNULL(mr.fee_charges_writtenoff_derived,0)) as fee_charges_writtenoff_derived,"); + updateSqlBuilder.append("SUM(IFNULL(mr.penalty_charges_amount,0)) as penalty_charges_charged_derived,"); + updateSqlBuilder.append("SUM(IFNULL(mr.penalty_charges_completed_derived,0)) as penalty_charges_repaid_derived,"); + updateSqlBuilder.append("SUM(IFNULL(mr.penalty_charges_waived_derived,0)) as penalty_charges_waived_derived,"); + updateSqlBuilder.append("SUM(IFNULL(mr.penalty_charges_writtenoff_derived,0)) as penalty_charges_writtenoff_derived "); + updateSqlBuilder.append(" FROM m_loan ml "); + updateSqlBuilder.append("INNER JOIN m_loan_repayment_schedule mr on mr.loan_id = ml.id "); + updateSqlBuilder.append("INNER JOIN m_client mc on mc.id=ml.client_id "); + updateSqlBuilder.append("INNER JOIN m_office o on mc.office_id = o.id "); + updateSqlBuilder.append("WHERE ml.disbursedon_date is not null and o.hierarchy like ? and ml.id < ? "); + updateSqlBuilder.append("GROUP BY ml.id "); + updateSqlBuilder.append("limit ? "); + updateSqlBuilder.append(") x on x.loanId = m_loan.id "); + + updateSqlBuilder.append("SET m_loan.principal_disbursed_derived = x.principal_disbursed_derived,"); + updateSqlBuilder.append("m_loan.principal_repaid_derived = x.principal_repaid_derived,"); + updateSqlBuilder.append("m_loan.principal_writtenoff_derived = x.principal_writtenoff_derived,"); + updateSqlBuilder + .append("m_loan.principal_outstanding_derived = (x.principal_disbursed_derived - (x.principal_repaid_derived + x.principal_writtenoff_derived)),"); + updateSqlBuilder.append("m_loan.interest_charged_derived = x.interest_charged_derived,"); + updateSqlBuilder.append("m_loan.interest_repaid_derived = x.interest_repaid_derived,"); + updateSqlBuilder.append("m_loan.interest_waived_derived = x.interest_waived_derived,"); + updateSqlBuilder.append("m_loan.interest_writtenoff_derived = x.interest_writtenoff_derived,"); + updateSqlBuilder + .append("m_loan.interest_outstanding_derived = (x.interest_charged_derived - (x.interest_repaid_derived + x.interest_waived_derived + x.interest_writtenoff_derived)),"); + updateSqlBuilder.append("m_loan.fee_charges_charged_derived = x.fee_charges_charged_derived,"); + updateSqlBuilder.append("m_loan.fee_charges_repaid_derived = x.fee_charges_repaid_derived,"); + updateSqlBuilder.append("m_loan.fee_charges_waived_derived = x.fee_charges_waived_derived,"); + updateSqlBuilder.append("m_loan.fee_charges_writtenoff_derived = x.fee_charges_writtenoff_derived,"); + updateSqlBuilder + .append("m_loan.fee_charges_outstanding_derived = (x.fee_charges_charged_derived - (x.fee_charges_repaid_derived + x.fee_charges_waived_derived + x.fee_charges_writtenoff_derived)),"); + updateSqlBuilder.append("m_loan.penalty_charges_charged_derived = x.penalty_charges_charged_derived,"); + updateSqlBuilder.append("m_loan.penalty_charges_repaid_derived = x.penalty_charges_repaid_derived,"); + updateSqlBuilder.append("m_loan.penalty_charges_waived_derived = x.penalty_charges_waived_derived,"); + updateSqlBuilder.append("m_loan.penalty_charges_writtenoff_derived = x.penalty_charges_writtenoff_derived,"); + updateSqlBuilder + .append("m_loan.penalty_charges_outstanding_derived = (x.penalty_charges_charged_derived - (x.penalty_charges_repaid_derived + x.penalty_charges_waived_derived + x.penalty_charges_writtenoff_derived)),"); + updateSqlBuilder + .append("m_loan.total_expected_repayment_derived = (x.principal_disbursed_derived + x.interest_charged_derived + x.fee_charges_charged_derived + x.penalty_charges_charged_derived),"); + updateSqlBuilder + .append("m_loan.total_repayment_derived = (x.principal_repaid_derived + x.interest_repaid_derived + x.fee_charges_repaid_derived + x.penalty_charges_repaid_derived),"); + updateSqlBuilder + .append("m_loan.total_expected_costofloan_derived = (x.interest_charged_derived + x.fee_charges_charged_derived + x.penalty_charges_charged_derived),"); + updateSqlBuilder + .append("m_loan.total_costofloan_derived = (x.interest_repaid_derived + x.fee_charges_repaid_derived + x.penalty_charges_repaid_derived),"); + updateSqlBuilder + .append("m_loan.total_waived_derived = (x.interest_waived_derived + x.fee_charges_waived_derived + x.penalty_charges_waived_derived),"); + updateSqlBuilder + .append("m_loan.total_writtenoff_derived = (x.interest_writtenoff_derived + x.fee_charges_writtenoff_derived + x.penalty_charges_writtenoff_derived),"); + updateSqlBuilder.append("m_loan.total_outstanding_derived="); + updateSqlBuilder.append(" (x.principal_disbursed_derived - (x.principal_repaid_derived + x.principal_writtenoff_derived)) + "); + updateSqlBuilder + .append(" (x.interest_charged_derived - (x.interest_repaid_derived + x.interest_waived_derived + x.interest_writtenoff_derived)) +"); + updateSqlBuilder + .append(" (x.fee_charges_charged_derived - (x.fee_charges_repaid_derived + x.fee_charges_waived_derived + x.fee_charges_writtenoff_derived)) +"); + updateSqlBuilder + .append(" (x.penalty_charges_charged_derived - (x.penalty_charges_repaid_derived + x.penalty_charges_waived_derived + x.penalty_charges_writtenoff_derived))"); + + final int result = jdbcTemplate.update(updateSqlBuilder.toString(), new Object[] {officeHierachy,lastLoanId,limit}); + + logger.info(ThreadLocalContextUtil.getTenant().getName() + ": Results affected by update: " + result); + + } +} \ No newline at end of file diff --git a/fineract-provider/src/main/java/org/apache/fineract/portfolio/loanaccount/service/LoanReadPlatformService.java b/fineract-provider/src/main/java/org/apache/fineract/portfolio/loanaccount/service/LoanReadPlatformService.java index e205fbf19ed..6ca3c71cdfd 100755 --- a/fineract-provider/src/main/java/org/apache/fineract/portfolio/loanaccount/service/LoanReadPlatformService.java +++ b/fineract-provider/src/main/java/org/apache/fineract/portfolio/loanaccount/service/LoanReadPlatformService.java @@ -134,4 +134,6 @@ LoanScheduleData retrieveRepaymentSchedule(Long loanId, RepaymentScheduleRelated LoanAccountData retrieveLoanByLoanAccount(String loanAccountNumber); Long retrieveLoanIdByAccountNumber(String loanAccountNumber); + + Integer retrieveNumberOfActiveLoans(); } \ No newline at end of file diff --git a/fineract-provider/src/main/java/org/apache/fineract/portfolio/loanaccount/service/LoanReadPlatformServiceImpl.java b/fineract-provider/src/main/java/org/apache/fineract/portfolio/loanaccount/service/LoanReadPlatformServiceImpl.java index e2464d83bbd..8087626346b 100755 --- a/fineract-provider/src/main/java/org/apache/fineract/portfolio/loanaccount/service/LoanReadPlatformServiceImpl.java +++ b/fineract-provider/src/main/java/org/apache/fineract/portfolio/loanaccount/service/LoanReadPlatformServiceImpl.java @@ -2209,4 +2209,9 @@ public Long retrieveLoanIdByAccountNumber(String loanAccountNumber) { } } + @Override + public Integer retrieveNumberOfActiveLoans() { + final String sql="select count(*) from m_loan"; + return this.jdbcTemplate.queryForObject(sql,Integer.class); + } } diff --git a/fineract-provider/src/main/java/org/apache/fineract/scheduledjobs/service/ScheduledJobRunnerService.java b/fineract-provider/src/main/java/org/apache/fineract/scheduledjobs/service/ScheduledJobRunnerService.java index 5199b751435..b6306b6c942 100644 --- a/fineract-provider/src/main/java/org/apache/fineract/scheduledjobs/service/ScheduledJobRunnerService.java +++ b/fineract-provider/src/main/java/org/apache/fineract/scheduledjobs/service/ScheduledJobRunnerService.java @@ -41,4 +41,6 @@ public interface ScheduledJobRunnerService { void postDividends() throws JobExecutionException; void updateTrialBalanceDetails() throws JobExecutionException; + + void updateLoanSummaryDetails(@SuppressWarnings("unused") final Map jobParameters) throws JobExecutionException; } diff --git a/fineract-provider/src/main/java/org/apache/fineract/scheduledjobs/service/ScheduledJobRunnerServiceImpl.java b/fineract-provider/src/main/java/org/apache/fineract/scheduledjobs/service/ScheduledJobRunnerServiceImpl.java index 29a337cd67d..abeb792c5f4 100644 --- a/fineract-provider/src/main/java/org/apache/fineract/scheduledjobs/service/ScheduledJobRunnerServiceImpl.java +++ b/fineract-provider/src/main/java/org/apache/fineract/scheduledjobs/service/ScheduledJobRunnerServiceImpl.java @@ -21,14 +21,17 @@ import java.math.BigDecimal; import java.math.BigInteger; import java.text.SimpleDateFormat; -import java.util.Collection; -import java.util.Date; -import java.util.List; -import java.util.Map; +import java.util.*; +import java.util.concurrent.*; import org.apache.fineract.accounting.glaccount.domain.TrialBalance; import org.apache.fineract.accounting.glaccount.domain.TrialBalanceRepositoryWrapper; +import org.apache.fineract.organisation.office.domain.Office; +import org.apache.fineract.organisation.office.domain.OfficeRepository; import org.apache.fineract.portfolio.loanaccount.api.LoanApiConstants; +import org.apache.fineract.portfolio.loanaccount.service.LoanApplicationWriteService; +import org.apache.fineract.portfolio.loanaccount.service.LoanReadPlatformService; +import org.apache.fineract.portfolio.loanaccount.service.UpdateLoanSummaryPoster; import org.joda.time.LocalDate; import org.joda.time.DateTime; @@ -57,6 +60,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.ApplicationContext; import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; @@ -77,6 +81,10 @@ public class ScheduledJobRunnerServiceImpl implements ScheduledJobRunnerService private final ShareAccountDividendReadPlatformService shareAccountDividendReadPlatformService; private final ShareAccountSchedularService shareAccountSchedularService; private final TrialBalanceRepositoryWrapper trialBalanceRepositoryWrapper; + private final LoanReadPlatformService loanReadPlatformService; + private final ApplicationContext applicationContext; + private final LoanApplicationWriteService loanApplicationWriteService; + private final OfficeRepository officeRepository; @Autowired public ScheduledJobRunnerServiceImpl(final RoutingDataSourceServiceFactory dataSourceServiceFactory, @@ -85,7 +93,11 @@ public ScheduledJobRunnerServiceImpl(final RoutingDataSourceServiceFactory dataS final DepositAccountReadPlatformService depositAccountReadPlatformService, final DepositAccountWritePlatformService depositAccountWritePlatformService, final ShareAccountDividendReadPlatformService shareAccountDividendReadPlatformService, - final ShareAccountSchedularService shareAccountSchedularService, final TrialBalanceRepositoryWrapper trialBalanceRepositoryWrapper) { + final ShareAccountSchedularService shareAccountSchedularService, final TrialBalanceRepositoryWrapper trialBalanceRepositoryWrapper, + final LoanReadPlatformService loanReadPlatformService, + final ApplicationContext applicationContext, + final LoanApplicationWriteService loanApplicationWriteService, + final OfficeRepository officeRepository) { this.dataSourceServiceFactory = dataSourceServiceFactory; this.savingsAccountWritePlatformService = savingsAccountWritePlatformService; this.savingsAccountChargeReadPlatformService = savingsAccountChargeReadPlatformService; @@ -94,6 +106,10 @@ public ScheduledJobRunnerServiceImpl(final RoutingDataSourceServiceFactory dataS this.shareAccountDividendReadPlatformService = shareAccountDividendReadPlatformService; this.shareAccountSchedularService = shareAccountSchedularService; this.trialBalanceRepositoryWrapper=trialBalanceRepositoryWrapper; + this.loanReadPlatformService=loanReadPlatformService; + this.applicationContext=applicationContext; + this.loanApplicationWriteService=loanApplicationWriteService; + this.officeRepository=officeRepository; } @Transactional @@ -483,4 +499,67 @@ public void updateTrialBalanceDetails() throws JobExecutionException { } + @Transactional + @Override + @CronTarget(jobName = JobName.UPDATE_LOAN_SUMMARY) + public void updateLoanSummaryDetails(@SuppressWarnings("unused") final Map jobParameters) throws JobExecutionException { + + final int threadPoolSize=Integer.parseInt(jobParameters.get("thread-pool-size")); + final String officeId = jobParameters.get("officeId"); + final ExecutorService executorService = Executors.newFixedThreadPool(threadPoolSize); + final Office office = this.officeRepository.findOne(Long.parseLong(officeId)); + + List> posters = new ArrayList>(); + + //Get the total count of loans + Integer numberofLoans=loanReadPlatformService.retrieveNumberOfActiveLoans(); + + Integer batchSize = (int) Math.ceil(numberofLoans/ threadPoolSize); + Integer maxLoanId=0; + + do{ + UpdateLoanSummaryPoster poster = (UpdateLoanSummaryPoster) this.applicationContext.getBean("updateLoanSummaryPoster"); + poster.setLoanApplicationWriteService(loanApplicationWriteService); + poster.setTenant(ThreadLocalContextUtil.getTenant()); + poster.setBatchSize(batchSize); + maxLoanId+=batchSize+1; + poster.setMaxLoanId(maxLoanId); + poster.setOfficeHierachy(office.getHierarchy()+ "%"); + posters.add(Executors.callable(poster)); + numberofLoans-=batchSize; + }while(numberofLoans>=0); + + try { + List> responses = executorService.invokeAll(posters); + checkCompletion(responses); + } catch (InterruptedException e1) { + logger.error("Interrupted while updating loan summary details", e1); + } + } + + //checks the execution of task by each thread in the executor service + private void checkCompletion(List> responses) { + try { + for(Future f : responses) { + f.get(); + } + boolean allThreadsExecuted = false; + int noOfThreadsExecuted = 0; + for (Future future : responses) { + if (future.isDone()) { + noOfThreadsExecuted++; + } + } + allThreadsExecuted = noOfThreadsExecuted == responses.size(); + if(!allThreadsExecuted) + logger.error("All threads could not execute."); + } catch (InterruptedException e1) { + logger.error("Interrupted while posting IRD entries", e1); + } catch (ExecutionException e2) { + logger.error("Execution exception while posting IRD entries", e2); + } + + } + + }