diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java index 208158ed2..d5fb99007 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java @@ -15,6 +15,7 @@ import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdf01Request; import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdfRequestPart; import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.clearing.ContinueSessionBnRequest; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.service.executors.AbstractExecutor; @@ -23,11 +24,9 @@ import ru.spcex.platform.classes.base.SpcexObjectBase; import ru.spcex.platform.enumeration.SdfTable; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.utils.collection.Pair; -import java.util.Collection; -import java.util.EnumMap; -import java.util.List; -import java.util.Map; +import java.util.*; import java.util.stream.Collectors; @Service @@ -39,6 +38,11 @@ public class StatementService extends QueueConsumer implements InitializingBean private final Map> sdfImdgs; private final Map> executorsMap; + /** + * Пары sdf запросов пришедшие с модуля dbf-import + */ + private final List> pairOfSdfRequest = new LinkedList<>(); + @Autowired public StatementService(Consumer kafkaQueue, ImdgProvider imdgProvider, @@ -66,17 +70,50 @@ public class StatementService extends QueueConsumer implements InitializingBean Collection sdfGroup; SdfTable table = statementRequest.getTable(); Imdg sdfImdg = sdfImdgs.get(table); + + Optional> completePairOpt = saveRequest(statementRequest); + if (completePairOpt.isPresent()) { + if (List.of(SdfTable.SDF_01, SdfTable.SDF_57).contains(table)) { + { + //всегда сначала обработаем sdf57 + Pair pair = completePairOpt.get(); + processSdf57(pair.getSecond()); + + //затем sdf01 + processSdf01(pair.getFirst()); + //теперь можем продолжить сессию с шага 1 + ContinueSessionBnRequest continueSessionBn = new ContinueSessionBnRequest(); + kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN_FIRST_PART, continueSessionBn); + } + } else if (List.of(SdfTable.SDF_08, SdfTable.SDF_13).contains(table)) { + //not implemented part; it's actually stage number 8 from any session + } + } + } + + private void processSdf57(StatementRequest statementRequest) { + Imdg sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf57, SDf57.class); + Collection sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of( + "generationId", statementRequest.getGroupId())); + AbstractExecutor service = executorsMap.get(SdfTable.SDF_57); + service.execute(sdfGroup, statementRequest); + } + + private void processSdf01(StatementRequest statementRequest) { + Imdg sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf01, SDf01.class); + Collection sdfGroup; if (statementRequest.getAccountCreationResults().size() == 0) { - sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of("generationId", statementRequest.getGroupId())); + sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of( + "generationId", statementRequest.getGroupId())); } else { + //todo могут ли тут реквесты бегать по кругу, если да, то у нас проблемы sdfGroup = statementRequest.getAccountCreationResults() .stream() .filter(part -> part.getErrorCode() == null) //fixme эти случае должны попадать в ошибочный sdf02 .map(part -> sdfImdg.getSingleObjectByID(part.getSdfId())) .collect(Collectors.toList()); } - - AbstractExecutor service = executorsMap.get(table); + AbstractExecutor service = executorsMap.get(SdfTable.SDF_01); Result res = service.execute(sdfGroup, statementRequest); if (res.getAccountRequests().size() != 0) { kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF01, createAccountsRequest(statementRequest.getGroupId(), res.getAccountRequests())); @@ -85,6 +122,50 @@ public class StatementService extends QueueConsumer implements InitializingBean } } + private Optional> saveRequest(StatementRequest statementRequest) { + SdfTable sdfTable = statementRequest.getTable(); + if (pairOfSdfRequest.isEmpty()) { + pairOfSdfRequest.add(getPairByTableName(sdfTable, statementRequest)); + } else { + //найдем первую неполноценную пару + Optional> uncompletedPairOpt = findUncompletedPairByTableName(sdfTable); + if (uncompletedPairOpt.isPresent()) { + Pair uncompletedPair = uncompletedPairOpt.get(); + if (uncompletedPair.getFirst() == null) { + uncompletedPair.setFirst(statementRequest); + } else { + uncompletedPair.setSecond(statementRequest); + } + //укомплектованная пара + return Optional.of(uncompletedPair); + } else { + //если таких нет, то просто создаем новую с одной частью + pairOfSdfRequest.add(getPairByTableName(sdfTable, statementRequest)); + } + } + return Optional.empty(); + } + + private Pair getPairByTableName(SdfTable sdfTable, StatementRequest statementRequest) { + Pair pair = null; + switch (sdfTable) { + case SDF_01 -> pair = new Pair<>(statementRequest, null); + case SDF_57 -> pair = new Pair<>(null, statementRequest); + } + return pair; + } + + private Optional> findUncompletedPairByTableName(SdfTable sdfTable) { + Optional> uncompletedPair = Optional.empty(); + switch (sdfTable) { + case SDF_01 -> + uncompletedPair = pairOfSdfRequest.stream().filter(pair -> pair.getFirst() != null && pair.getSecond() == null).findFirst(); + case SDF_57 -> + uncompletedPair = pairOfSdfRequest.stream().filter(pair -> pair.getSecond() != null && pair.getFirst() == null).findFirst(); + } + return uncompletedPair; + } + private AccountSdf01Request createAccountsRequest(Long sdf01GroupingId, List accountRequests) { AccountSdf01Request r = new AccountSdf01Request(); r.setGroupingSdf01Id(sdf01GroupingId); diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SdfTable.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SdfTable.java index 0d885df19..6a920d430 100644 --- a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SdfTable.java +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SdfTable.java @@ -3,7 +3,9 @@ package ru.spcex.platform.enumeration; import ru.spcex.platform.utils.enumeration.IEnumKey; public enum SdfTable implements IEnumKey { - SDF_01("SDF_01"), SDF_57("SDF_57"), SDF_09("SDF_09"), SDF_16("SDF_16"); + SDF_01("SDF_01"), SDF_04("SDF_04"), SDF_08("SDF_08"), + SDF_13("SDF_13"), SDF_57("SDF_57"), SDF_09("SDF_09"), + SDF_16("SDF_16"); SdfTable(String key) { this.key = key;