From f9fa84ad64f7dd875143b0205d7428ed1422cad3 Mon Sep 17 00:00:00 2001 From: AKurakin Date: Fri, 9 Jun 2023 16:50:12 +0300 Subject: [PATCH] =?UTF-8?q?clearing-service=20sdf08=20=D0=BE=D0=B1=D1=80?= =?UTF-8?q?=D0=B0=D0=B1=D0=BE=D1=82=D0=BA=D0=B0=20=D0=B0=D0=B2=D1=82=D0=BE?= =?UTF-8?q?=D1=81=D0=BE=D0=B7=D0=B4=D0=B0=D0=BD=D0=B8=D1=8F=20=D1=81=D1=87?= =?UTF-8?q?=D0=B5=D1=82=D0=BE=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../spcex/clearing/service/StatementService.java | 16 ++++++++++++---- .../logic/stages/DbfImportKafkaMessenger.java | 9 +++++++-- 2 files changed, 19 insertions(+), 6 deletions(-) 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 d8ec4bde6..c2a748287 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 @@ -123,7 +123,8 @@ public class StatementService extends QueueConsumer implements InitializingBean Collection sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of( "generationId", statementRequest.getGroupId())); AbstractExecutor service = executorsMap.get(SdfTable.SDF_57); - service.execute(sdfGroup, statementRequest); + Result res = service.execute(sdfGroup, statementRequest); + finishSendCommand(res, service, statementRequest); } private void processSdf04(StatementRequest statementRequest) { @@ -132,7 +133,8 @@ public class StatementService extends QueueConsumer implements InitializingBean "generationId", statementRequest.getGroupId())); AbstractExecutor service = executorsMap.get(SdfTable.SDF_04); if (service != null) { - service.execute(sdfGroup, statementRequest); + Result res = service.execute(sdfGroup, statementRequest); + finishSendCommand(res, service, statementRequest); } else { log.warn("Executor for SDF_04 not set"); } @@ -144,7 +146,8 @@ public class StatementService extends QueueConsumer implements InitializingBean "generationId", statementRequest.getGroupId())); AbstractExecutor service = executorsMap.get(SdfTable.SDF_08); if (service != null) { - service.execute(sdfGroup, statementRequest); + Result res = service.execute(sdfGroup, statementRequest); + finishSendCommand(res, service, statementRequest); } else { log.warn("Executor for SDF_08 not set"); } @@ -155,7 +158,8 @@ public class StatementService extends QueueConsumer implements InitializingBean Collection sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of( "generationId", statementRequest.getGroupId())); AbstractExecutor service = executorsMap.get(SdfTable.SDF_13); - service.execute(sdfGroup, statementRequest); + Result res = service.execute(sdfGroup, statementRequest); + finishSendCommand(res, service, statementRequest); } private void processSdf01(StatementRequest statementRequest) { @@ -174,6 +178,10 @@ public class StatementService extends QueueConsumer implements InitializingBean } AbstractExecutor service = executorsMap.get(SdfTable.SDF_01); Result res = service.execute(sdfGroup, statementRequest); + finishSendCommand(res, service, statementRequest); + } + + void finishSendCommand(Result res, AbstractExecutor service, StatementRequest statementRequest) { if (res.getAccountRequests().size() != 0) { kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF01, createAccountsRequest(statementRequest.getGroupId(), res.getAccountRequests())); } else if (service.isNeedToSendCommand()) { 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 ae14aecb4..eed6b0b63 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 @@ -1,5 +1,7 @@ package ru.spcex.clearing.dbf.importer.logic.stages; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; import org.springframework.stereotype.Component; import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable; @@ -16,6 +18,7 @@ import java.util.function.Supplier; @Component public class DbfImportKafkaMessenger implements InitializingBean { + final Logger log = LoggerFactory.getLogger(getClass()); private final Supplier kafka; private final Map> messengers; @@ -30,7 +33,7 @@ public class DbfImportKafkaMessenger implements InitializingBean { messengers.put(ETable.DF_09, groupId -> messageBalance(groupId, SdfTable.SDF_09)); messengers.put(ETable.DF_16, groupId -> messageBalance(groupId, SdfTable.SDF_16)); messengers.put(ETable.DF_57, groupId -> messageBalance(groupId, SdfTable.SDF_57)); - messengers.put(ETable.DF_04, groupId -> messageBalance(groupId, SdfTable.SDF_04)); + messengers.put(ETable.DF_04, groupId -> messageBalance(groupId, SdfTable.SDF_04)); } /** @@ -48,7 +51,9 @@ public class DbfImportKafkaMessenger implements InitializingBean { StatementRequest statementRequest = new StatementRequest(); statementRequest.setGroupId(groupId); statementRequest.setTable(table); - kafka.get().sendRequestToQueue(Consts.STATEMENT_PROCESS, statementRequest); + Long msgId = kafka.get().sendRequestToQueue(Consts.STATEMENT_PROCESS, statementRequest); + log.debug("Send StatementRequest({}, {}) message id={} to kafka \"{}\"", + groupId, table, msgId, Consts.STATEMENT_PROCESS); } private void messageDf04(Long groupId) {