diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/ConsumerSpecificClass.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/ConsumerSpecificClass.java index 6d47624d2..597cf8b04 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/ConsumerSpecificClass.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/ConsumerSpecificClass.java @@ -5,26 +5,24 @@ import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import java.util.function.BiConsumer; import java.util.function.Consumer; +/** + * вспомогательный класс + * callback для пришедшего сообщения из очереди + */ public class ConsumerSpecificClass implements BuilderConsumerStep, BuilderDestinationStep { + /** + * какой payload будет в BaseRequest + */ private Class clazz; + /** + * действие-обработчик реквеста (метод из конкретного менеджера) + */ private java.util.function.Consumer> consumer; - public ConsumerSpecificClass(Class clazz) { + private ConsumerSpecificClass(Class clazz) { this.clazz = clazz; } - public void forDestination(String destination, BiConsumer> fun) { - fun.accept(destination, this); - } - - public void acceptRaw(Object obj) { - this.consumer.accept((BaseRequest) obj); - } - - public Class getClazz() { - return clazz; - } - public static BuilderConsumerStep build(Class clazz) { return new ConsumerSpecificClass<>(clazz); } @@ -34,4 +32,17 @@ public class ConsumerSpecificClass implements BuilderConsumerStep, Builder this.consumer = consumer; return this; } + + public void forDestination(String destination, BiConsumer> fun) { + fun.accept(destination, this); + } + + @SuppressWarnings("unchecked") + public void acceptRaw(Object obj) { + this.consumer.accept((BaseRequest) obj); + } + + public Class getClazz() { + return clazz; + } } 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 6918f3262..2bcc35950 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 @@ -21,6 +21,9 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.atomic.AtomicBoolean; +/** + * утилитный класс для обработки сообщений из очереди + */ public class QueueConsumer implements AutoCloseable { private final Logger log = LoggerFactory.getLogger(getClass()); private final AtomicBoolean closed = new AtomicBoolean(false); @@ -69,7 +72,7 @@ public class QueueConsumer implements AutoCloseable { } @Override - public void close() throws Exception { + public void close() { closed.set(true); consumer.wakeup(); }