From ce767ceae121be27612c7603abebb98e9001da49 Mon Sep 17 00:00:00 2001 From: ialbert Date: Tue, 29 Nov 2022 19:21:56 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-130 clearing-service --- clearing-parent/clearing-service/pom.xml | 72 +++++++++ .../clearing/ClearingServiceApplication.java | 12 ++ .../clearing/config/ClearingImdgConfig.java | 47 ++++++ .../ru/spcex/clearing/config/KafkaConfig.java | 49 +++++++ .../config/MessageResolverConfig.java | 18 +++ .../spcex/clearing/config/MessagesConfig.java | 29 ++++ .../element/ClearingServiceSettings.java | 41 ++++++ .../spcex/clearing/error/ClearingError.java | 21 +++ .../clearing/service/ClearingService.java | 77 ++++++++++ .../clearing/service/PaymentBatchInfo.java | 46 ++++++ .../service/PaymentBatchProcessor.java | 138 ++++++++++++++++++ .../src/main/resources/application.properties | 22 +++ clearing-parent/pom.xml | 1 + .../platform/enumeration/AccountType.java | 2 +- .../enumeration/ClearingMemberCategoryD.java | 18 +++ .../enumeration/TransactionStatus.java | 18 +++ 16 files changed, 610 insertions(+), 1 deletion(-) create mode 100644 clearing-parent/clearing-service/pom.xml create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/ClearingServiceApplication.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ClearingImdgConfig.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/KafkaConfig.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MessageResolverConfig.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MessagesConfig.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/element/ClearingServiceSettings.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/error/ClearingError.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ClearingService.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentBatchInfo.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentBatchProcessor.java create mode 100644 clearing-parent/clearing-service/src/main/resources/application.properties create mode 100644 platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/ClearingMemberCategoryD.java create mode 100644 platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/TransactionStatus.java diff --git a/clearing-parent/clearing-service/pom.xml b/clearing-parent/clearing-service/pom.xml new file mode 100644 index 000000000..2c40979e9 --- /dev/null +++ b/clearing-parent/clearing-service/pom.xml @@ -0,0 +1,72 @@ + + + + clearing-parent + ru.spcex.clearing + SPCEX-1.0.0.0 + + 4.0.0 + + clearing-service + + + 17 + 17 + + + + + ru.spcex.platform + platform-messaging + + + ru.spcex.platform + platform-imdg-api-hazelcast-impl + + + ru.spcex.clearing + classes + + + org.springframework.boot + spring-boot-starter + + + com.fasterxml.jackson.core + jackson-databind + + + ru.spcex.platform + platform-enum + + + + + + src/main/resources + + application.properties + + false + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + repackage + + + + + ${project.artifactId} + + + + + \ No newline at end of file diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/ClearingServiceApplication.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/ClearingServiceApplication.java new file mode 100644 index 000000000..3a4c724f8 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/ClearingServiceApplication.java @@ -0,0 +1,12 @@ +package ru.spcex.clearing; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +@SpringBootApplication +public class ClearingServiceApplication { + public static void main(String[] args) { + SpringApplication app = new SpringApplication(ClearingServiceApplication.class); + app.run(args); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ClearingImdgConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ClearingImdgConfig.java new file mode 100644 index 000000000..670a536a7 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ClearingImdgConfig.java @@ -0,0 +1,47 @@ +package ru.spcex.clearing.config; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; +import ru.spcex.clearing.config.element.ClearingServiceSettings; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; + +@Configuration +public class ClearingImdgConfig { + private static ThreadPoolTaskExecutor createThreadPoolTaskExecutor(int maxPoolSz, boolean waitForCompletion) { + ThreadPoolTaskExecutor pool = new ThreadPoolTaskExecutor(); + if (maxPoolSz > 2) { + pool.setKeepAliveSeconds(60); + pool.setAllowCoreThreadTimeOut(true); + } + pool.setCorePoolSize(maxPoolSz); + pool.setWaitForTasksToCompleteOnShutdown(waitForCompletion); + return pool; + } + + @Bean(name = "taskExecutorHazelcastClientInitializer") + public ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer() { + return createThreadPoolTaskExecutor(1, true); + } + + @Bean(name = "taskExecutorIdGeneratorAwaiter") + public ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter() { + return createThreadPoolTaskExecutor(1, false); + } + + @Autowired + @Bean + public ImdgProvider imdgProvider( + @Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer, + @Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter, + ClearingServiceSettings settings + ) { + return new HazelcastService(taskExecutorHazelcastClientInitializer, + taskExecutorIdGeneratorAwaiter, + settings.getHazelcast()); + } + +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/KafkaConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/KafkaConfig.java new file mode 100644 index 000000000..0a281b5e1 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/KafkaConfig.java @@ -0,0 +1,49 @@ +package ru.spcex.clearing.config; + +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.config.ConfigurableBeanFactory; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Scope; +import ru.spcex.clearing.config.element.ClearingServiceSettings; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; +import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; +import ru.spcex.clearing.platform.messaging.service.RequestInfo; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgId; +import ru.spcex.platform.imdg.api.ImdgProvider; + +@Configuration +public class KafkaConfig { + @Autowired + @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) + @Bean + public Consumer createConsumer(ClearingServiceSettings settings) { + return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); + } + + @Autowired + @Bean + public Producer createProducer(ClearingServiceSettings settings) { + return KafkaProducerFactory.producer(settings.getKafkaProducer()); + } + + @Autowired + @Bean + public KafkaSender kafkaSender(Producer kafkaProducer, ImdgProvider imdgProvider) { + ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); + return KafkaSender + .setup() + .producer(kafkaProducer) + .idGenerator(imdgIdGenerator::nextId) + .imdgProvider(s -> { + Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); + return imdg::insert; + }) + .build(); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MessageResolverConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MessageResolverConfig.java new file mode 100644 index 000000000..09157e14d --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MessageResolverConfig.java @@ -0,0 +1,18 @@ +package ru.spcex.clearing.config; + +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import ru.spcex.platform.utils.enumeration.IMessageResolver; + +import java.util.Arrays; + +@Configuration +public class MessageResolverConfig { + @Bean + public IMessageResolver messageResolver() { + return errorMessage -> { + if (errorMessage == null) return "null"; + return String.format("(%d) args %s", errorMessage.getSubject().getId(), Arrays.toString(errorMessage.getArgs())); + }; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MessagesConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MessagesConfig.java new file mode 100644 index 000000000..38d52587e --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MessagesConfig.java @@ -0,0 +1,29 @@ +package ru.spcex.clearing.config; + +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.context.annotation.Bean; +import org.springframework.context.support.ResourceBundleMessageSource; +import ru.spcex.platform.utils.enumeration.IMessageResolver; +import ru.spcex.platform.utils.enumeration.SpringPropertiesMessageResolver; + +import java.util.Locale; + +//@Configuration +public class MessagesConfig { + @Bean("validation-error-messages") + public ResourceBundleMessageSource messages() { + ResourceBundleMessageSource source = new ResourceBundleMessageSource(); + source.setBasenames("messages/error"); + source.setUseCodeAsDefaultMessage(true); + source.setDefaultEncoding("utf8"); + source.setDefaultLocale(Locale.ROOT); + return source; + } + + @Bean + public IMessageResolver errorResolver(@Qualifier("validation-error-messages") ResourceBundleMessageSource messageBundle) { + SpringPropertiesMessageResolver resolver = new SpringPropertiesMessageResolver(messageBundle); + resolver.setLocale("ru"); + return resolver; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/element/ClearingServiceSettings.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/element/ClearingServiceSettings.java new file mode 100644 index 000000000..d42e5fcc3 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/element/ClearingServiceSettings.java @@ -0,0 +1,41 @@ +package ru.spcex.clearing.config.element; + +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.context.annotation.PropertySource; +import org.springframework.stereotype.Component; +import ru.spcex.clearing.platform.messaging.config.element.KafkaConsumerSettings; +import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings; +import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams; + +@Component +@PropertySource("file:${spring.config.location}/application.properties") +@ConfigurationProperties("clearing-service") +public class ClearingServiceSettings { + private HazelcastClientParams hazelcast; + private KafkaConsumerSettings kafkaConsumer; + private KafkaProducerSettings kafkaProducer; + + public HazelcastClientParams getHazelcast() { + return hazelcast; + } + + public void setHazelcast(HazelcastClientParams hazelcast) { + this.hazelcast = hazelcast; + } + + public KafkaConsumerSettings getKafkaConsumer() { + return kafkaConsumer; + } + + public void setKafkaConsumer(KafkaConsumerSettings kafkaConsumer) { + this.kafkaConsumer = kafkaConsumer; + } + + public KafkaProducerSettings getKafkaProducer() { + return kafkaProducer; + } + + public void setKafkaProducer(KafkaProducerSettings kafkaProducer) { + this.kafkaProducer = kafkaProducer; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/error/ClearingError.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/error/ClearingError.java new file mode 100644 index 000000000..9e30eb31a --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/error/ClearingError.java @@ -0,0 +1,21 @@ +package ru.spcex.clearing.error; + +import ru.spcex.platform.utils.enumeration.IEnumId; + +public enum ClearingError implements IEnumId { + CompanyCreditCheck(10012L), + CompanyDebitCheck(10013L), + //if ever happens, ask to add + ClearingMemberCategoryUnknown(19999L), + ; + private final Long id; + + ClearingError(Long id) { + this.id = id; + } + + @Override + public Long getId() { + return id; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ClearingService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ClearingService.java new file mode 100644 index 000000000..64eabf711 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ClearingService.java @@ -0,0 +1,77 @@ +package ru.spcex.clearing.service; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.scheduling.annotation.EnableScheduling; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.payment.PaymentInstruction; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.platform.enumeration.TransactionStatus; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgId; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.util.AbstractMap; +import java.util.Comparator; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.stream.Collectors; + +@Service +@EnableScheduling +public class ClearingService { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg paymentImdgs; + private final ExecutorService executor; + private final ImdgId idGenerator; + private final PaymentBatchProcessor senderGroupProcessor; + + @Autowired + public ClearingService(ImdgProvider imdgProvider, PaymentBatchProcessor senderGroupProcessor) { + this.paymentImdgs = imdgProvider.getImdg(IMDGDistributedNames.Map_PaymentInstruction, PaymentInstruction.class); + this.idGenerator = imdgProvider.getImdgIdGenerator(); + this.senderGroupProcessor = senderGroupProcessor; + this.executor = Executors.newSingleThreadExecutor(); + } + + @Scheduled(cron = "${clearing-service.scheduler.check-payment-instruction}") + public void run() { + executor.execute(this::createSdfFromPaymentInstructionSTLD); + } + + private void createSdfFromPaymentInstructionSTLD() { + //выгружаем PaymentInstructions с нужным статусом, группируем по компаниям + Map> pmtInsBySender = paymentImdgs.getCollectionObjectsByFieldValues( + Map.of("transactionStatus", TransactionStatus.stld.getKey())) + .stream() + .sorted(Comparator.comparing(PaymentInstruction::getSenderId)) + .collect(Collectors.groupingBy(PaymentInstruction::getSenderId)); + //generationId для созадаваемых Sdf03/Sdf11 + Long generationId = idGenerator.nextId(); + //результатом работы senderGroupProcessor будет Map PaymentBatchInfo> + //PaymentBatchInfo содержит возможную ошибку, при необходимости отсортированные Payment + //тип ClearingMemberCategory + Map processResults = pmtInsBySender + .entrySet() + .stream() + .map(entry -> { + Long senderId = entry.getKey(); + List pmtInstrcs = entry.getValue(); + return new AbstractMap.SimpleEntry<>(senderId, senderGroupProcessor.processSingleCompanyPayments(generationId, senderId, pmtInstrcs)); + }) + .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); + + //todo export + //set status + + } + + private void processSingleCompanyPayments(Long generationId, Long senderId, List payments) { + log.info("processing PaymentInstruction's generationId={} senderId={} size={}", generationId, senderId, payments.size()); + + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentBatchInfo.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentBatchInfo.java new file mode 100644 index 000000000..77ba4ad1a --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentBatchInfo.java @@ -0,0 +1,46 @@ +package ru.spcex.clearing.service; + +import ru.clearing.classes.statics.data.payment.PaymentInstruction; +import ru.spcex.platform.enumeration.ClearingMemberCategoryD; +import ru.spcex.platform.utils.enumeration.EnumMessage; + +import java.util.List; + +public class PaymentBatchInfo { + private List fromClearingToBank; + private List fromBankToClearing; + private ClearingMemberCategoryD categoryD; + private EnumMessage error; + + public List getFromClearingToBank() { + return fromClearingToBank; + } + + public void setFromClearingToBank(List fromClearingToBank) { + this.fromClearingToBank = fromClearingToBank; + } + + public List getFromBankToClearing() { + return fromBankToClearing; + } + + public void setFromBankToClearing(List fromBankToClearing) { + this.fromBankToClearing = fromBankToClearing; + } + + public ClearingMemberCategoryD getCategoryD() { + return categoryD; + } + + public void setCategoryD(ClearingMemberCategoryD categoryD) { + this.categoryD = categoryD; + } + + public EnumMessage getError() { + return error; + } + + public void setError(EnumMessage error) { + this.error = error; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentBatchProcessor.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentBatchProcessor.java new file mode 100644 index 000000000..bdd1b5dcb --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentBatchProcessor.java @@ -0,0 +1,138 @@ +package ru.spcex.clearing.service; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Component; +import ru.clearing.classes.statics.data.account.Account; +import ru.clearing.classes.statics.data.generated.ClearingMemberCategory; +import ru.clearing.classes.statics.data.payment.PaymentInstruction; +import ru.spcex.clearing.error.ClearingError; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.platform.enumeration.AccountType; +import ru.spcex.platform.enumeration.ClearingMemberCategoryD; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.utils.enumeration.EnumMessage; +import ru.spcex.platform.utils.enumeration.IEnumKey; +import ru.spcex.platform.utils.enumeration.IMessageResolver; + +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.function.Function; + +@Component +public class PaymentBatchProcessor { + + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg clrngMmbrImdg; + private final Imdg accImdg; + private final IMessageResolver messageResolver; + + @Autowired + public PaymentBatchProcessor(ImdgProvider imdgProvider, IMessageResolver messageResolver) { + this.clrngMmbrImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingMemberCategory, ClearingMemberCategory.class); + this.accImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class); + + this.messageResolver = messageResolver; + } + + PaymentBatchInfo processSingleCompanyPayments(Long generationId, Long senderId, List payments) { + PaymentBatchInfo batchInfo = new PaymentBatchInfo(); + log.info("processing PaymentInstruction's generationId={} senderId={} size={}", generationId, senderId, payments.size()); + ClearingMemberCategory category = clrngMmbrImdg.getSingleObjectByFieldValues(Map.of("companyId", senderId)); + ClearingMemberCategoryD categoryValue; + if (category == null || (categoryValue = + IEnumKey.getEnumByKey(ClearingMemberCategoryD.class, category.getClearingMemberCategory())) == null) { + log.warn("processing PaymentInstructions generationId={} senderId={}: " + + "ClearingMemberCategory {}, ClearingMemberCategory.clearingMemberCategory {}", + generationId, senderId, + category == null ? "null" : "is not null", + category == null ? "null" : category.getClearingMemberCategory() + ); + batchInfo.setError(new EnumMessage(ClearingError.ClearingMemberCategoryUnknown)); + return batchInfo; + } + batchInfo.setCategoryD(categoryValue); + if (ClearingMemberCategoryD.I.equals(categoryValue)) { + List aList = new ArrayList<>(); + List bList = new ArrayList<>(); + for (PaymentInstruction payment : payments) { + Account creditLegAcc = accImdg.getSingleObjectByID(payment.getCreditLegAccountId()); + Account debitLegAcc = accImdg.getSingleObjectByID(payment.getDebitLegAccountId()); + if (creditLegAcc == null || debitLegAcc == null) { + log.warn("generationId={}, senderId={} payment.id={} cannot find account CreditLegAccountId/DebitLegAccountId {}/{}", + generationId, senderId, payment.getId(), payment.getCreditLegAccountId(), payment.getDebitLegAccountId()); + continue; + } + if (AccountType.Clrn.equalsByKey(creditLegAcc.getAccountType()) + && AccountType.Bank.equalsByKey(debitLegAcc.getAccountType())) { + aList.add(payment); + } else if (AccountType.Bank.equalsByKey(creditLegAcc.getAccountType()) + && AccountType.Clrn.equalsByKey(debitLegAcc.getAccountType())) { + bList.add(payment); + } else { + log.warn("generationId={} senderId={} cannot sort payment.id={}, creditLegAccount.type={}, debitLegAccount.type={}", + generationId, senderId, payment.getId(),creditLegAcc.getAccountType(), debitLegAcc.getAccountType()); + continue; + } + } + Function, Long> creditAmountSum = paymentInstructions -> paymentInstructions + .stream() + .map(PaymentInstruction::getCreditLegAmount) + .reduce(0L, Long::sum); + Function, Long> debitAmountSum = paymentInstructions -> paymentInstructions + .stream() + .map(PaymentInstruction::getCreditLegAmount) + .reduce(0L, Long::sum); + Long fromClearingToBankCreditAmount = creditAmountSum.apply(aList); + Long fromBankToClearingCreditAmount = creditAmountSum.apply(bList); + Long fromClearingToBankDebitAmount = debitAmountSum.apply(aList); + Long fromBankToClearingDebitAmount = debitAmountSum.apply(bList); + if (!fromClearingToBankCreditAmount.equals(fromBankToClearingCreditAmount)) { + log.info("processing PaymentInstruction's generationId={} senderId={} [fromClearingToBankCreditAmount={}, " + + "fromBankToClearingCreditAmount={}, " + + "fromClearingToBankDebitAmount={}, " + + "fromBankToClearingDebitAmount={}] error {}", generationId, senderId, + fromClearingToBankCreditAmount, + fromBankToClearingCreditAmount, + fromClearingToBankDebitAmount, + fromBankToClearingDebitAmount, + messageResolver.resolve(new EnumMessage(ClearingError.CompanyCreditCheck, senderId.toString()))); + batchInfo.setError(new EnumMessage(ClearingError.CompanyCreditCheck, senderId.toString())); + return batchInfo; + } else if (!fromClearingToBankDebitAmount.equals(fromBankToClearingDebitAmount)) { + log.info("processing PaymentInstruction's generationId={} senderId={} [fromClearingToBankCreditAmount={}, " + + "fromBankToClearingCreditAmount={}, " + + "fromClearingToBankDebitAmount={}, " + + "fromBankToClearingDebitAmount={}] error {}", generationId, senderId, + fromClearingToBankCreditAmount, + fromBankToClearingCreditAmount, + fromClearingToBankDebitAmount, + fromBankToClearingDebitAmount, + messageResolver.resolve(new EnumMessage(ClearingError.CompanyDebitCheck, senderId.toString()))); + batchInfo.setError(new EnumMessage(ClearingError.CompanyDebitCheck, senderId.toString())); + return batchInfo; + } + log.debug("processing PaymentInstruction's generationId={} senderId={} [fromClearingToBankCreditAmount={}, " + + "fromBankToClearingCreditAmount={}, " + + "fromClearingToBankDebitAmount={}, " + + "fromBankToClearingDebitAmount={}]", generationId, senderId, + fromClearingToBankCreditAmount, + fromBankToClearingCreditAmount, + fromClearingToBankDebitAmount, + fromBankToClearingDebitAmount); + batchInfo.setFromClearingToBank(aList); + batchInfo.setFromBankToClearing(bList); + return batchInfo; + } else if (ClearingMemberCategoryD.B.equals(categoryValue)) { + return batchInfo; + } else { + log.warn("PaymentInstructions generationId={} senderId={} fail - unknown category ", + generationId, senderId); + batchInfo.setError(new EnumMessage(ClearingError.ClearingMemberCategoryUnknown)); + return batchInfo; + } + } +} diff --git a/clearing-parent/clearing-service/src/main/resources/application.properties b/clearing-parent/clearing-service/src/main/resources/application.properties new file mode 100644 index 000000000..af4ebcf34 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/resources/application.properties @@ -0,0 +1,22 @@ +spring.main.web-application-type=none + +clearing-service.hazelcast.cluster-members=127.0.0.1:5701 +clearing-service.hazelcast.login=dev +clearing-service.hazelcast.password=dev-pass + +clearing-service.kafka-consumer.bootstrap-servers=localhost:9092 +clearing-service.kafka-consumer.group-id=dev-group-clearing-service +clearing-service.kafka-consumer.enable-auto-commit=false +clearing-service.kafka-consumer.session-timeout-ms=30000 +clearing-service.kafka-consumer.auto-offset-reset=latest +clearing-service.kafka-consumer.linger-ms=1 +clearing-service.kafka-consumer.buffer-memory=33554432 + +clearing-service.kafka-producer.bootstrap-servers=localhost:9092 +clearing-service.kafka-producer.acks=all +clearing-service.kafka-producer.retries=0 +clearing-service.kafka-producer.batch-size=16384 +clearing-service.kafka-producer.linger-ms=1 +clearing-service.kafka-producer.buffer-memory=33554432 + +clearing-service.scheduler.check-payment-instruction=*/5 * * * * * \ No newline at end of file diff --git a/clearing-parent/pom.xml b/clearing-parent/pom.xml index 08cc8557b..85fb2a0a6 100644 --- a/clearing-parent/pom.xml +++ b/clearing-parent/pom.xml @@ -30,6 +30,7 @@ account-service balance-service scheduler-service + clearing-service diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/AccountType.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/AccountType.java index 71385a733..bf0f5f50b 100644 --- a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/AccountType.java +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/AccountType.java @@ -3,7 +3,7 @@ package ru.spcex.platform.enumeration; import ru.spcex.platform.utils.enumeration.IEnumKey; public enum AccountType implements IEnumKey { - Clrn("CLRN"); + Clrn("CLRN"), Bank("BANK"); private final String key; diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/ClearingMemberCategoryD.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/ClearingMemberCategoryD.java new file mode 100644 index 000000000..bc94d111c --- /dev/null +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/ClearingMemberCategoryD.java @@ -0,0 +1,18 @@ +package ru.spcex.platform.enumeration; + +import ru.spcex.platform.utils.enumeration.IEnumKey; + +public enum ClearingMemberCategoryD implements IEnumKey { + I("I"), B("B"); + + private final String key; + + ClearingMemberCategoryD(String key) { + this.key = key; + } + + @Override + public String getKey() { + return key; + } +} diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/TransactionStatus.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/TransactionStatus.java new file mode 100644 index 000000000..812d075bd --- /dev/null +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/TransactionStatus.java @@ -0,0 +1,18 @@ +package ru.spcex.platform.enumeration; + +import ru.spcex.platform.utils.enumeration.IEnumKey; + +public enum TransactionStatus implements IEnumKey { + stld("STLD"); + + private final String key; + + TransactionStatus(String key) { + this.key = key; + } + + @Override + public String getKey() { + return key; + } +}