This commit is contained in:
parent
a127b85e12
commit
394bb7477c
2 changed files with 28 additions and 14 deletions
|
|
@ -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<T> implements BuilderConsumerStep<T>, BuilderDestinationStep {
|
||||
/**
|
||||
* какой payload будет в BaseRequest
|
||||
*/
|
||||
private Class<T> clazz;
|
||||
/**
|
||||
* действие-обработчик реквеста (метод из конкретного менеджера)
|
||||
*/
|
||||
private java.util.function.Consumer<BaseRequest<T>> consumer;
|
||||
|
||||
public ConsumerSpecificClass(Class<T> clazz) {
|
||||
private ConsumerSpecificClass(Class<T> clazz) {
|
||||
this.clazz = clazz;
|
||||
}
|
||||
|
||||
public void forDestination(String destination, BiConsumer<String, ConsumerSpecificClass<?>> fun) {
|
||||
fun.accept(destination, this);
|
||||
}
|
||||
|
||||
public void acceptRaw(Object obj) {
|
||||
this.consumer.accept((BaseRequest<T>) obj);
|
||||
}
|
||||
|
||||
public Class<T> getClazz() {
|
||||
return clazz;
|
||||
}
|
||||
|
||||
public static <T1> BuilderConsumerStep<T1> build(Class<T1> clazz) {
|
||||
return new ConsumerSpecificClass<>(clazz);
|
||||
}
|
||||
|
|
@ -34,4 +32,17 @@ public class ConsumerSpecificClass<T> implements BuilderConsumerStep<T>, Builder
|
|||
this.consumer = consumer;
|
||||
return this;
|
||||
}
|
||||
|
||||
public void forDestination(String destination, BiConsumer<String, ConsumerSpecificClass<?>> fun) {
|
||||
fun.accept(destination, this);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
public void acceptRaw(Object obj) {
|
||||
this.consumer.accept((BaseRequest<T>) obj);
|
||||
}
|
||||
|
||||
public Class<T> getClazz() {
|
||||
return clazz;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue