From 7ff3de272424703fbc2f2e8ee7a222dee252174c Mon Sep 17 00:00:00 2001 From: AKurakin Date: Wed, 5 Apr 2023 09:36:56 +0300 Subject: [PATCH] =?UTF-8?q?CLS-259=20util=20=D0=B4=D0=BE=D0=B1=D0=B0=D0=B2?= =?UTF-8?q?=D0=B8=D0=BB=20=D0=B0=D1=81=D1=81=D0=B8=D0=BD=D1=85=D1=80=D0=BE?= =?UTF-8?q?=D0=BD=D0=BD=D0=BE-=D1=81=D0=B8=D0=BD=D1=85=D1=80=D0=BE=D0=BD?= =?UTF-8?q?=D0=BD=D1=8B=D0=B9=20=D0=BE=D0=B1=D0=BC=D0=B5=D0=BD=20=D1=81?= =?UTF-8?q?=D0=BE=D0=BE=D0=B1=D1=89=D0=B5=D0=BD=D0=B8=D1=8F=D0=BC=D0=B8=20?= =?UTF-8?q?(=D0=B8=D0=BD=D1=82=D0=B5=D1=80=D1=84=D0=B5=D0=B9=D1=81)=20BiDi?= =?UTF-8?q?rectionQueueExchanger?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../exchangers/BiDirectionQueueExchanger.java | 66 +++++++++++++++++++ 1 file changed, 66 insertions(+) create mode 100644 clearing-parent/clearing-utils/src/main/java/ru/spcex/clearing/util/services/exchangers/BiDirectionQueueExchanger.java diff --git a/clearing-parent/clearing-utils/src/main/java/ru/spcex/clearing/util/services/exchangers/BiDirectionQueueExchanger.java b/clearing-parent/clearing-utils/src/main/java/ru/spcex/clearing/util/services/exchangers/BiDirectionQueueExchanger.java new file mode 100644 index 000000000..0f08fe624 --- /dev/null +++ b/clearing-parent/clearing-utils/src/main/java/ru/spcex/clearing/util/services/exchangers/BiDirectionQueueExchanger.java @@ -0,0 +1,66 @@ +package ru.spcex.clearing.util.services.exchangers; + +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.DisposableBean; +import org.springframework.beans.factory.InitializingBean; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; + +/** + * Синхронный обмен сообщениями с ассинхронным сервисом. + * Метод exchange() отправляет сообщение в очередь и дожидается ответа из другой очереди + */ +public class BiDirectionQueueExchanger> extends QueueConsumer implements InitializingBean, DisposableBean { + protected final Logger log = LoggerFactory.getLogger(getClass()); + protected final Object sync = new Object(); + protected String outQueue; + protected String inQueue; + protected Class listenClass; + protected long timeout; + + /** + * Синхронно-ассинхронный обмен сообщениями + * + * @param kafkaQueue + * @param kafkaProducer + * @param outQueue отправляет в очередь + * @param inQueue слушает очередь, ожидает ответов + * @param listenClass типы объектов из inQueue + * @param timeout - максимальное ожидание ответа, в миллисекундах, 0 - неограничено + */ + public BiDirectionQueueExchanger(Consumer kafkaQueue, Producer kafkaProducer, + String outQueue, + String inQueue, Class listenClass, + long timeout) { + super(kafkaQueue, kafkaProducer); + this.outQueue = outQueue; + this.inQueue = inQueue; + this.listenClass = listenClass; + this.timeout = timeout; + } + + /** + * Отправить сообщение message в outQueue и дождаться ответа из очереди inQueue + * + * @param message + * @return + * @throws InterruptedException + */ + public TOut exchange(TIn message) throws InterruptedException { + //todo impl BiDirectionQueueExchanger + return null; + } + + @Override + public void destroy() throws Exception { + log.debug("Listener {} for async exchange {}-{} close", this, outQueue, inQueue); + } + + @Override + public void afterPropertiesSet() throws Exception { + log.debug("Listener {} for async exchange {}-{} ready", this, outQueue, inQueue); + } +}