From 864bbcf125a12fbd40bf9421b4cf1c9c09cca18f Mon Sep 17 00:00:00 2001 From: AKurakin Date: Thu, 4 May 2023 16:29:52 +0300 Subject: [PATCH] =?UTF-8?q?company-service=20http://jira.mfd.msk:8088/brow?= =?UTF-8?q?se/CLS-272=20=D0=B8=D1=81=D0=BF=D1=80=D0=B0=D0=B2=D0=B8=D0=BB?= =?UTF-8?q?=20=D0=B8=D1=81=D0=BF=D0=BE=D0=BB=D1=8C=D0=B7=D0=BE=D0=B2=D0=B0?= =?UTF-8?q?=D0=BD=D0=B8=D0=B5=20BiDirectionQueueExchanger?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../util/services/exchangers/BiDirectionQueueExchanger.java | 2 +- .../spcex/clearing/company/service/ClientCodeService.java | 6 +++--- .../clearing/company/service/ClientCodeServiceTest.java | 1 - .../clearing/platform/messaging/service/QueueConsumer.java | 1 + 4 files changed, 5 insertions(+), 5 deletions(-) diff --git a/clearing-parent/clearing-validation/src/main/java/ru/spcex/clearing/util/services/exchangers/BiDirectionQueueExchanger.java b/clearing-parent/clearing-validation/src/main/java/ru/spcex/clearing/util/services/exchangers/BiDirectionQueueExchanger.java index c26f63b0c..b90b1993e 100644 --- a/clearing-parent/clearing-validation/src/main/java/ru/spcex/clearing/util/services/exchangers/BiDirectionQueueExchanger.java +++ b/clearing-parent/clearing-validation/src/main/java/ru/spcex/clearing/util/services/exchangers/BiDirectionQueueExchanger.java @@ -55,7 +55,7 @@ public class BiDirectionQueueExchanger> extends Queu /** * Синхронно-ассинхронный обмен сообщениями * - * @param kafkaQueue + * @param kafkaQueue уникальный kafka Consumer (Spring prototype). Нельзя переиспользовать существующие. * @param kafkaProducer * @param outQueue отправляет в очередь * @param inQueue слушает очередь/топик, ожидает ответов. Пример: Consts.CONTINUE_CLEARING diff --git a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/ClientCodeService.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/ClientCodeService.java index 2e82f5e9b..a795a26db 100644 --- a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/ClientCodeService.java +++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/ClientCodeService.java @@ -70,7 +70,7 @@ public class ClientCodeService extends QueueConsumer implements InitializingBean BiDirectionQueueExchanger> accountServiceExchanger; @Autowired - public ClientCodeService(Consumer kafkaQueue, + public ClientCodeService(Consumer kafkaQueue1, Consumer kafkaQueue2, Producer kafkaProducer, ImdgProvider imdgProvider, ValidationHelper validationHelper, @@ -79,7 +79,7 @@ public class ClientCodeService extends QueueConsumer implements InitializingBean @Qualifier("clientCodeNewRequestValidator") Function clientCodeNewRequestValidator, @Qualifier("clientCodeUpdateRequestValidator") Function clientCodeUpdateRequestValidator, @Qualifier("clientCodeDeleteRequestValidator") Function clientCodeDeleteRequestValidator) { - super(kafkaQueue, kafkaProducer); + super(kafkaQueue1, kafkaProducer); this.kafkaProducer = kafkaProducer; this.idGenerator = imdgProvider.getImdgIdGenerator(); this.clientCodeMap = imdgProvider.getImdg(IMDGDistributedNames.Map_ClientCode, ClientCode.class); @@ -92,7 +92,7 @@ public class ClientCodeService extends QueueConsumer implements InitializingBean this.messageResolver = messageResolver; this.userRoleVerification = userRoleVerification; - this.accountServiceExchanger = new BiDirectionQueueExchanger<>(kafkaQueue, kafkaProducer, + this.accountServiceExchanger = new BiDirectionQueueExchanger<>(kafkaQueue2, kafkaProducer, Consts.DESTINATION_TRADING_CLEARING_REGISTRY_NEW, Consts.REQUEST_INFO_UPDATE, true, // REQUEST_INFO_UPDATE - стандартная очередь, для результатов всех реквестов. DESTINATION_TRADING_CLEARING_REGISTRY_REPLY, 60000 diff --git a/clearing-parent/company-service/src/test/java/ru/spcex/clearing/company/service/ClientCodeServiceTest.java b/clearing-parent/company-service/src/test/java/ru/spcex/clearing/company/service/ClientCodeServiceTest.java index 503174d22..0f34d0778 100644 --- a/clearing-parent/company-service/src/test/java/ru/spcex/clearing/company/service/ClientCodeServiceTest.java +++ b/clearing-parent/company-service/src/test/java/ru/spcex/clearing/company/service/ClientCodeServiceTest.java @@ -45,7 +45,6 @@ import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparato import static ru.spcex.clearing.test.TestUtils.*; import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId; -//todo тесты случайно ломаются из-за BiDirectionQueueExchanger. Когда там выключаю initReplyListener() то тесты стабильно работают. @ExtendWith(SpringExtension.class) @ContextConfiguration(classes = { ClientCodeService.class, diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java index 0b26d5c03..3f7afe713 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java @@ -74,6 +74,7 @@ public class QueueConsumer implements AutoCloseable { public void init() { inputExecutor.submit(() -> { try { + log.debug("Using kafka consumer {} for subscribe on \"{}\"", consumer, callbacks.keySet()); if (supportStartOffsetTimeWindow) { consumer.subscribe(callbacks.keySet(), new OffsetChanger(consumer, callbacks.keySet())); } else {