From d28b12f68e94d00f547e042b0f0631f16b231f93 Mon Sep 17 00:00:00 2001 From: ialbert Date: Mon, 12 Dec 2022 17:08:35 +0300 Subject: [PATCH 1/6] http://jira.mfd.msk:8088/browse/CLS-36 bug fix --- .../importer/logic/stages/DbfImportKafkaMessenger.java | 2 ++ .../dbf/importer/services/DBFImporterService.java | 6 ++++-- .../imdg/iml/hazelcast/adapter/ImdgHazelcast.java | 8 ++++---- 3 files changed, 10 insertions(+), 6 deletions(-) diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java index e72cf67f4..2706c2476 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java @@ -7,6 +7,7 @@ import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest; import ru.spcex.clearing.platform.messaging.domain.cud.clearing.Sdf04Request; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.enumeration.SdfTable; import java.util.HashMap; import java.util.Map; @@ -43,6 +44,7 @@ public class DbfImportKafkaMessenger implements InitializingBean { private void messageDf01(Long groupId) { StatementRequest statementRequest = new StatementRequest(); statementRequest.setGroupId(groupId); + statementRequest.setTable(SdfTable.SDF_01); kafka.get().sendRequestToQueue(Consts.STATEMENT_PROCESS, statementRequest); } diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java index 51eed3b59..fd3b2cdce 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java @@ -45,10 +45,12 @@ public class DBFImporterService { executorService.execute(() -> { log.info("checking new files... {}", specificTable != null ? specificTable.name() : ""); Map> newFiles = fileChecker.checkNewFiles(specificTable); - if (log.isDebugEnabled()) { + if (newFiles.size() > 0 && log.isDebugEnabled()) { log.debug("following files will be processed {}", forLogging(newFiles)); - } else { + } else if (newFiles.size() > 0) { log.info("following files will be processed {}", forLoggingSizeOnly(newFiles)); + } else { + log.info("no files were found"); } for (Map.Entry> newFilesEntry : newFiles.entrySet()) { ETable currTable = newFilesEntry.getKey(); diff --git a/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/adapter/ImdgHazelcast.java b/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/adapter/ImdgHazelcast.java index 0db5e0b28..b1e431044 100644 --- a/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/adapter/ImdgHazelcast.java +++ b/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/adapter/ImdgHazelcast.java @@ -90,8 +90,8 @@ public class ImdgHazelcast implements Imdg { predicates[i[0]] = Predicates.equal(key, value); i[0]++; }); - Predicate or = Predicates.or(predicates); - Set> found = map.entrySet(or); + Predicate and = Predicates.and(predicates); + Set> found = map.entrySet(and); Iterator> allFoundByCondition = found.iterator(); if (allFoundByCondition.hasNext()) { return allFoundByCondition.next().getValue(); @@ -122,8 +122,8 @@ public class ImdgHazelcast implements Imdg { predicates[i[0]] = Predicates.equal(key, value); i[0]++; }); - Predicate or = Predicates.or(predicates); - Set ids = map.keySet(or); + Predicate and = Predicates.and(predicates); + Set ids = map.keySet(and); Iterator idIterator = ids.iterator(); Collection searchResult = new ArrayList<>(); while (idIterator.hasNext()) { From 3bc6c83b1c797ea5f9eb76c790e609422aafdaf1 Mon Sep 17 00:00:00 2001 From: ialbert Date: Mon, 12 Dec 2022 18:54:33 +0300 Subject: [PATCH 2/6] http://jira.mfd.msk:8088/browse/CLS-36 parallel file execution temp fix --- .../importer/services/DBFImporterService.java | 63 +++++++++++++------ 1 file changed, 45 insertions(+), 18 deletions(-) diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java index fd3b2cdce..62a911533 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java @@ -3,7 +3,6 @@ package ru.spcex.clearing.dbf.importer.services; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Qualifier; -import org.springframework.scheduling.annotation.EnableScheduling; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; import org.springframework.stereotype.Service; @@ -12,9 +11,8 @@ import ru.spcex.clearing.dbf.importer.logic.data.ResultContainer; import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable; import java.io.File; -import java.util.LinkedList; -import java.util.List; -import java.util.Map; +import java.nio.file.Path; +import java.util.*; import java.util.function.Function; import java.util.stream.Collectors; import java.util.stream.Stream; @@ -26,6 +24,7 @@ public class DBFImporterService { private final FileChecker fileChecker; private final ThreadPoolTaskExecutor executorService; private final Processor processor; + private final Set filesCurrentlyInProcess; public DBFImporterService(@Qualifier("fileChecker") FileChecker messageListener, @Qualifier("executor") ThreadPoolTaskExecutor executorService, @@ -33,6 +32,7 @@ public class DBFImporterService { this.fileChecker = messageListener; this.executorService = executorService; this.processor = processor; + this.filesCurrentlyInProcess = new HashSet<>(); } @Scheduled(cron = "${import-dbf-service.scheduler.check-src-dir-cron}") @@ -43,25 +43,52 @@ public class DBFImporterService { public void run(ETable specificTable) { log.info("adding import task {}", specificTable != null ? specificTable.name() : ""); executorService.execute(() -> { - log.info("checking new files... {}", specificTable != null ? specificTable.name() : ""); - Map> newFiles = fileChecker.checkNewFiles(specificTable); - if (newFiles.size() > 0 && log.isDebugEnabled()) { - log.debug("following files will be processed {}", forLogging(newFiles)); - } else if (newFiles.size() > 0) { - log.info("following files will be processed {}", forLoggingSizeOnly(newFiles)); - } else { - log.info("no files were found"); - } - for (Map.Entry> newFilesEntry : newFiles.entrySet()) { - ETable currTable = newFilesEntry.getKey(); - List fileList = newFilesEntry.getValue(); - for (File dbfFile : fileList) { - processor.process(ResultContainer.createNewTask(currTable, dbfFile)); + Map> newFiles = null; + try { + log.info("checking new files... {}", specificTable != null ? specificTable.name() : ""); + newFiles = getFiles(specificTable); + if (newFiles.size() > 0 && log.isDebugEnabled()) { + log.debug("following files will be processed {}", forLogging(newFiles)); + } else if (newFiles.size() > 0) { + log.info("following files will be processed {}", forLoggingSizeOnly(newFiles)); + } else { + log.info("no files were found"); + } + for (Map.Entry> newFilesEntry : newFiles.entrySet()) { + ETable currTable = newFilesEntry.getKey(); + List fileList = newFilesEntry.getValue(); + for (File dbfFile : fileList) { + processor.process(ResultContainer.createNewTask(currTable, dbfFile)); + } + } + } finally { + if (newFiles != null && newFiles.size() > 0) { + cleanFiles(newFiles); } } }); } + private synchronized Map> getFiles(ETable specificTable) { + Map> newFiles = fileChecker.checkNewFiles(specificTable); + Map> newFilesFiltered = newFiles.entrySet() + .stream() + .map(entry -> + new AbstractMap.SimpleEntry<>(entry.getKey(), entry.getValue() + .stream() + .filter(file -> !filesCurrentlyInProcess.contains(file.toPath())) + .collect(Collectors.toList()))) + .collect(Collectors.toMap(AbstractMap.SimpleEntry::getKey, AbstractMap.SimpleEntry::getValue)); + newFilesFiltered.forEach((table, files) + -> files.forEach(file -> filesCurrentlyInProcess.add(file.toPath()))); + return newFilesFiltered; + } + + private synchronized void cleanFiles(Map> filesFromTask) { + filesFromTask.forEach((table, files) + -> files.forEach(file -> filesCurrentlyInProcess.remove(file.toPath()))); + } + private static String forLogging(Map> files) { return files .entrySet() From 48f02b2af3ec3f837ab99b39ee80f0e47c5b3b63 Mon Sep 17 00:00:00 2001 From: ialbert Date: Mon, 12 Dec 2022 18:56:10 +0300 Subject: [PATCH 3/6] http://jira.mfd.msk:8088/browse/CLS-36 --- .../spcex/clearing/dbf/importer/services/DBFImporterService.java | 1 + 1 file changed, 1 insertion(+) diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java index 62a911533..a2f4e64b6 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java @@ -3,6 +3,7 @@ package ru.spcex.clearing.dbf.importer.services; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.scheduling.annotation.EnableScheduling; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; import org.springframework.stereotype.Service; From 33193b97464855f6feab4f3aee700187a38fce6f Mon Sep 17 00:00:00 2001 From: ialbert Date: Tue, 13 Dec 2022 14:22:47 +0300 Subject: [PATCH 4/6] http://jira.mfd.msk:8088/browse/CLS-42 --- .../balance/service/Sdf08Service.java | 36 +++++++--- .../ru/spcex/platform/enumeration/Task.java | 2 +- .../domain/cud/balance/SDf08NewRequest.java | 67 ------------------- 3 files changed, 26 insertions(+), 79 deletions(-) delete mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/balance/SDf08NewRequest.java diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf08Service.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf08Service.java index b7b2143c4..d93076049 100644 --- a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf08Service.java +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf08Service.java @@ -9,39 +9,53 @@ import ru.clearing.classes.statics.data.sdf.SDf08; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.Consts; -import ru.spcex.clearing.platform.messaging.domain.cud.balance.SDf08NewRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.balance.ExportToFileRequest; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.enumeration.Task; import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgId; import ru.spcex.platform.imdg.api.ImdgProvider; +import java.math.BigDecimal; +import java.time.Instant; + @Service public class Sdf08Service extends QueueConsumer implements InitializingBean { private final Logger log = LoggerFactory.getLogger(getClass()); private final Imdg sdf8Map; + private final ImdgId idGenerator; + private final KafkaSender kafkaReqProducer; - public Sdf08Service(Consumer kafkaQueue, ImdgProvider imdgProvider) { + public Sdf08Service(Consumer kafkaQueue, ImdgProvider imdgProvider, KafkaSender kafkaReqProducer) { super(kafkaQueue); this.sdf8Map = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf08, SDf08.class); + this.idGenerator = imdgProvider.getImdgIdGenerator(); + this.kafkaReqProducer = kafkaReqProducer; } @Override public void afterPropertiesSet() { - callback(SDf08NewRequest.class) + callback(Object.class) .setConsumer(this::newSDf08) - .forDestination(Consts.DESTINATION_SDF08_NEW, callbacks::put); + .forDestination(Task.getAllBalance.topic(), callbacks::put); init(); } - private void newSDf08(BaseRequest userRequest) { - SDf08NewRequest req = userRequest.getRequestPayload(); - log.debug("SDf08NewRequest received"); + private void newSDf08(BaseRequest userRequest) { + log.debug("getAllBalance request received"); SDf08 sDf08 = new SDf08(); - sDf08.setNumber(req.getNumber()); - sDf08.setDatetime(req.getDatetime()); - sDf08.setGenerationTime(req.getGenerationTime()); - sDf08.setGenerationId(req.getGenerationId()); + sDf08.setNumber(BigDecimal.valueOf(Math.random())); + Instant now = Instant.now(); + sDf08.setDatetime(now); + sDf08.setGenerationTime(now); + sDf08.setGenerationId(idGenerator.nextId()); sdf8Map.insert(sDf08); + ExportToFileRequest exportRequest = new ExportToFileRequest(); + exportRequest.setSdfGroupId(sDf08.getGenerationId()); + exportRequest.setNameOfTable("DF-08"); + kafkaReqProducer.sendRequestToQueue(Consts.EXPORT_PROCESS, exportRequest); log.debug("successfully processed, new id {}", sDf08.getId()); } } diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java index 762af41b4..bf4bbb5fa 100644 --- a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java @@ -3,7 +3,7 @@ package ru.spcex.platform.enumeration; import ru.spcex.platform.utils.enumeration.IEnumKey; public enum Task implements IEnumKey { - createOrder("CORD"), createOrderConfirm("CORC"); + createOrder("CORD"), createOrderConfirm("CORC"), getAllBalance("GALB"); private final String key; diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/balance/SDf08NewRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/balance/SDf08NewRequest.java deleted file mode 100644 index b0038e7f9..000000000 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/balance/SDf08NewRequest.java +++ /dev/null @@ -1,67 +0,0 @@ -package ru.spcex.clearing.platform.messaging.domain.cud.balance; - -import com.fasterxml.jackson.annotation.JsonProperty; -import com.fasterxml.jackson.databind.annotation.JsonDeserialize; -import com.fasterxml.jackson.databind.annotation.JsonSerialize; -import ru.spcex.clearing.platform.messaging.domain.json.deserialize.InstantDeserializer; -import ru.spcex.clearing.platform.messaging.domain.json.serialize.InstantSerializer; - -import java.math.BigDecimal; -import java.time.Instant; - -public class SDf08NewRequest { - @JsonProperty - private BigDecimal number; - @JsonSerialize(using = InstantSerializer.class) - @JsonDeserialize(using = InstantDeserializer.class) - @JsonProperty - private Instant datetime; - @JsonProperty - private String fileName; - @JsonSerialize(using = InstantSerializer.class) - @JsonDeserialize(using = InstantDeserializer.class) - @JsonProperty - private Instant generationTime; - @JsonProperty - private Long generationId; - - public BigDecimal getNumber() { - return number; - } - - public void setNumber(BigDecimal number) { - this.number = number; - } - - public Instant getDatetime() { - return datetime; - } - - public void setDatetime(Instant datetime) { - this.datetime = datetime; - } - - public String getFileName() { - return fileName; - } - - public void setFileName(String fileName) { - this.fileName = fileName; - } - - public Instant getGenerationTime() { - return generationTime; - } - - public void setGenerationTime(Instant generationTime) { - this.generationTime = generationTime; - } - - public Long getGenerationId() { - return generationId; - } - - public void setGenerationId(Long generationId) { - this.generationId = generationId; - } -} From d7e570fa3925aa5cb65038619c881ebeb92c90de Mon Sep 17 00:00:00 2001 From: ialbert Date: Tue, 13 Dec 2022 15:15:30 +0300 Subject: [PATCH 5/6] http://jira.mfd.msk:8088/browse/CLS-205 --- .../MoneyMarketSecurityUpdateAction.java | 14 ++++++++++++++ .../MoneyMarketSecurityUpdateRequest.java | 13 +++++++++++++ 2 files changed, 27 insertions(+) diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/request/cud/securities/MoneyMarketSecurityUpdateAction.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/request/cud/securities/MoneyMarketSecurityUpdateAction.java index ae0124840..3928e8198 100644 --- a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/request/cud/securities/MoneyMarketSecurityUpdateAction.java +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/request/cud/securities/MoneyMarketSecurityUpdateAction.java @@ -16,6 +16,11 @@ public class MoneyMarketSecurityUpdateAction implements IAction Date: Tue, 13 Dec 2022 15:16:40 +0300 Subject: [PATCH 6/6] http://jira.mfd.msk:8088/browse/CLS-205 --- .../securities/service/cud/MoneyMarketSecurityService.java | 1 + 1 file changed, 1 insertion(+) diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java index c5d40b200..18029734e 100644 --- a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java @@ -112,6 +112,7 @@ public class MoneyMarketSecurityService extends QueueConsumer implements Initial log.debug("MoneyMarketSecurityUpdateRequest received id = {}", req.getId()); MoneyMarketSecurity mms = moneyMarketSecurityMap.getSingleObjectByID(req.getId()); mms.setUpdated(Instant.now()); + mms.setStartDate(req.getStartDate()); mms.setEndDate(req.getEndDate()); mms.setNominalValue(req.getNominalValue() != null ? BigDecimal.valueOf(req.getNominalValue()) : null); mms.setNominalCurrency(req.getNominalCurrency());