From 9f018d13c42da2970ab6e930a097bb5e1ed287a0 Mon Sep 17 00:00:00 2001 From: ialbert Date: Thu, 6 Jul 2023 19:03:15 +0300 Subject: [PATCH] Sdf08Executor registry capacity: DepoAccount search fix --- .../service/ClearingAccountService.java | 2 +- .../account/service/DepoAccountService.java | 2 +- .../clearing/service/StatementService.java | 98 +++++++++++++------ .../service/executors/Sdf08Executor.java | 3 +- .../domain/cud/balance/StatementRequest.java | 10 +- 5 files changed, 76 insertions(+), 39 deletions(-) diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java index e8b007e57..cd26a0eef 100644 --- a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java @@ -273,7 +273,7 @@ public class ClearingAccountService extends QueueConsumer implements Initializin StatementRequest request = new StatementRequest(); request.setGroupId(groupingSdf01Id); request.setAccountCreationResults(results); - request.setContinueSdf01(true); + request.setContinueSdf(true); request.setTable(SdfTable.SDF_01); // по нему запрос получили log.debug("Send message to kafka \"{}\": {}", Consts.STATEMENT_PROCESS, LogFormatter.toStringWrapper(request)); kafkaSender.sendRequestToQueue(Consts.STATEMENT_PROCESS, request); diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/DepoAccountService.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/DepoAccountService.java index 97a1375d6..710a0ed13 100644 --- a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/DepoAccountService.java +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/DepoAccountService.java @@ -211,7 +211,7 @@ public class DepoAccountService extends QueueConsumer implements InitializingBea public void sendStatementRequestBack(Long groupingSdf01Id, List results) { StatementRequest request = new StatementRequest(); request.setGroupId(groupingSdf01Id); - request.setContinueSdf01(true); //fixme???? + request.setContinueSdf(true); //fixme???? request.setAccountCreationResults(results); request.setTable(SdfTable.SDF_08); // по нему запрос получили log.debug("Send message to kafka \"{}\": {}", Consts.STATEMENT_PROCESS, LogFormatter.toStringWrapper(request)); 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 756a67ed7..ed22b4a91 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 @@ -66,20 +66,9 @@ public class StatementService extends QueueConsumer implements InitializingBean public void afterPropertiesSet() throws Exception { callback(StatementRequest.class) .setConsumer(systemRequest -> { - if (systemRequest.getRequestPayload() != null && systemRequest.getRequestPayload().isContinueSdf01()) { - Optional sdf01And57Key = findCompletePair(); - if (sdf01And57Key.isEmpty()) { - log.error("FATAL: response from account-service received id={} sdf ids={}", systemRequest.getId(), - systemRequest.getRequestPayload() - .getAccountCreationResults() - .stream() - .map(res -> String.valueOf(res.getSdfId())) - .collect(Collectors.joining(",", "[", "]")) - ); - return; - } - processSdf01And57(sdf01And57Key.get(), systemRequest.getRequestPayload()); - log.info("Sdf01 and Sdf57 processed successfully after account-service command."); + StatementRequest payload = systemRequest.getRequestPayload(); + if (payload.isContinueSdf()) { + processAccountAnswer(systemRequest); } else { processPaired(systemRequest); } @@ -88,10 +77,15 @@ public class StatementService extends QueueConsumer implements InitializingBean init(); } + /** + * fromAccService передается когда пришел ответ от account-service + * в этом случае: по key находим пару в которой сохранен sdf57 запрос и частично выполненный sdf01 + * вместо старого sdf01 запроса выполняем новый пришедший от account-service + */ private void processSdf01And57(Long key, StatementRequest fromAccService) { Pair pair = pairOfSdfRequest.get(key); - boolean allProcessed = processSdf01(fromAccService == null ? pair.getFirst() : fromAccService); - if (!allProcessed) { + Result sdf01Res = processSdf01(fromAccService == null ? pair.getFirst() : fromAccService); + if (sdf01Res.getAccountRequests().size() > 0) { log.info("sdf01 execution wasn't complete, waiting for an answer from account-service"); return; } @@ -106,6 +100,31 @@ public class StatementService extends QueueConsumer implements InitializingBean } + private void processAccountAnswer(BaseRequest systemRequest) { + StatementRequest payload = systemRequest.getRequestPayload(); + switch (payload.getTable()) { + case SDF_01 -> { + //там нужно переиспользовать SDF57 запрос приходивший ранее и сохраненный в pairOfSdfRequest + Optional key = findCompleteKey(payload.getTable()); + if (key.isEmpty()) { + log.error("FATAL: response from account-service received id={} table {} sdf ids={}", systemRequest.getId(), + payload.getTable(), + payload + .getAccountCreationResults() + .stream() + .map(res -> String.valueOf(res.getSdfId())) + .collect(Collectors.joining(",", "[", "]")) + ); + return; + } + processSdf01And57(key.get(), payload); + } + case SDF_08 -> processSdf08(payload); + } + log.info("{} processed successfully after account-service command.", payload.getTable()); + + } + private void processPaired(BaseRequest systemRequest) { StatementRequest statementRequest = systemRequest.getRequestPayload(); SdfTable table = statementRequest.getTable(); @@ -172,8 +191,17 @@ public class StatementService extends QueueConsumer implements InitializingBean private void processSdf08(StatementRequest statementRequest) { Imdg sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf08, SDf08.class); - Collection sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of( - "generationId", statementRequest.getGroupId())); + Collection sdfGroup; + if (statementRequest.getAccountCreationResults().size() == 0) { + sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of( + "generationId", statementRequest.getGroupId())); + } else { + sdfGroup = statementRequest.getAccountCreationResults() + .stream() + .filter(part -> part.getErrorCode() == null) + .map(part -> sdfImdg.getSingleObjectByID(part.getSdfId())) + .collect(Collectors.toList()); + } AbstractExecutor service = executorsMap.get(SdfTable.SDF_08); if (service != null) { Result res = service.execute(sdfGroup, statementRequest); @@ -196,7 +224,7 @@ public class StatementService extends QueueConsumer implements InitializingBean // finishSendCommand(res, service, statementRequest); } - private boolean processSdf01(StatementRequest statementRequest) { + private Result processSdf01(StatementRequest statementRequest) { Imdg sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf01, SDf01.class); Collection sdfGroup; if (statementRequest.getAccountCreationResults().size() == 0) { @@ -214,12 +242,12 @@ public class StatementService extends QueueConsumer implements InitializingBean Result res = service.execute(sdfGroup, statementRequest); if (res.getAccountRequests().size() != 0) { kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF01, createAccountsRequest(statementRequest.getGroupId(), res.getAccountRequests())); - return false; + return res; } else if (service.isNeedToSendCommand()) { service.sendCommand(kafkaSender, res); - return true; + return res; } - return true; + return res; } private Optional saveRequest(StatementRequest statementRequest) { @@ -282,14 +310,24 @@ public class StatementService extends QueueConsumer implements InitializingBean return uncompletedPair; } - private Optional findCompletePair() { - return pairOfSdfRequest.entrySet() - .stream() - .filter(entry -> entry.getValue().getFirst() != null && entry.getValue().getSecond() != null) - .filter(entry -> SdfTable.SDF_01.equals(entry.getValue().getFirst().getTable()) - && SdfTable.SDF_57.equals(entry.getValue().getSecond().getTable())) - .map(Map.Entry::getKey) - .findFirst(); + private Optional findCompleteKey(SdfTable sdfTable) { + Optional key = Optional.empty(); + switch (sdfTable) { + case SDF_01 -> key = pairOfSdfRequest.entrySet() + .stream() + .filter(entry -> entry.getValue().getFirst() != null && entry.getValue().getSecond() != null) + .filter(entry -> SdfTable.SDF_01.equals(entry.getValue().getFirst().getTable()) + && SdfTable.SDF_57.equals(entry.getValue().getSecond().getTable())) + .map(Map.Entry::getKey) + .findFirst(); + case SDF_08 -> key = pairOfSdfRequest.entrySet() + .stream() + .filter(entry -> entry.getValue().getFirst() != null) + .filter(entry -> SdfTable.SDF_08.equals(entry.getValue().getFirst().getTable())) + .map(Map.Entry::getKey) + .findFirst(); + } + return key; } private AccountSdf01Request createAccountsRequest(Long sdf01GroupingId, List accountRequests) { diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf08Executor.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf08Executor.java index 054c0e1d2..963aebea6 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf08Executor.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf08Executor.java @@ -42,7 +42,6 @@ import java.util.Collection; import java.util.Map; import java.util.Optional; import java.util.function.Function; -import java.util.stream.Collectors; @Service public class Sdf08Executor extends AbstractExecutor { @@ -267,7 +266,7 @@ public class Sdf08Executor extends AbstractExecutor { rgs.setAccount(account.getAccount()); rgs.setRegistryDesignation(RegistryDesignation.A.getKey()); rgs.setRegistryInstrumentType(RegistryInstrumentType.S.getKey()); - DepoAccount depoAccount = depoAccountImdg.getSingleObjectByID(statement.getAccountId()); + DepoAccount depoAccount = depoAccountImdg.getSingleObjectByFieldValues(Map.of("accountId", statement.getAccountId())); if (depoAccount != null) { rgs.setRegistryCapacity(depoAccount.getDepoAccountType()); } diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/balance/StatementRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/balance/StatementRequest.java index afa0a3981..eec240f6b 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/balance/StatementRequest.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/balance/StatementRequest.java @@ -16,18 +16,18 @@ public class StatementRequest { @JsonProperty List accountCreationResults = new ArrayList<>(); @JsonProperty - private boolean continueSdf01 = false; + private boolean continueSdf = false; public Long getGroupId() { return groupId; } - public boolean isContinueSdf01() { - return continueSdf01; + public boolean isContinueSdf() { + return continueSdf; } - public void setContinueSdf01(boolean continueSdf01) { - this.continueSdf01 = continueSdf01; + public void setContinueSdf(boolean continueSdf) { + this.continueSdf = continueSdf; } public void setGroupId(Long groupId) {