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 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/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..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 @@ -12,9 +12,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 +25,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 +33,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,23 +44,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 (log.isDebugEnabled()) { - log.debug("following files will be processed {}", forLogging(newFiles)); - } else { - log.info("following files will be processed {}", forLoggingSizeOnly(newFiles)); - } - 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() 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()); 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-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()) { 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; - } -} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/securitites/MoneyMarketSecurityUpdateRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/securitites/MoneyMarketSecurityUpdateRequest.java index 2ec1444e2..220e52c8a 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/securitites/MoneyMarketSecurityUpdateRequest.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/securitites/MoneyMarketSecurityUpdateRequest.java @@ -16,6 +16,11 @@ public class MoneyMarketSecurityUpdateRequest { @JsonSerialize(using = LocalDateSerializer.class) @JsonDeserialize(using = LocalDateDeserializer.class) @JsonProperty + public LocalDate startDate; + @JsonFormat(pattern = "yyyy-MM-dd", timezone = "Europe/Moscow") + @JsonSerialize(using = LocalDateSerializer.class) + @JsonDeserialize(using = LocalDateDeserializer.class) + @JsonProperty public LocalDate endDate; @JsonProperty public Double nominalValue; @@ -83,4 +88,12 @@ public class MoneyMarketSecurityUpdateRequest { public void setLotSize(Double lotSize) { this.lotSize = lotSize; } + + public LocalDate getStartDate() { + return startDate; + } + + public void setStartDate(LocalDate startDate) { + this.startDate = startDate; + } }