From ef1492822f12659b1435c3e0a1abdb5a75bdf30d Mon Sep 17 00:00:00 2001 From: ialbert Date: Thu, 11 Aug 2022 19:36:49 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-16 --- .../backendapi/config/MessagesConfig.java | 21 ++++ .../config/element/BackendApiSettings.java | 8 +- .../backendapi/meta/CudMetaService.java | 3 +- clearing-parent/pom.xml | 1 + clearing-parent/securities-service/pom.xml | 62 ++++++++++++ .../SecuritiesServiceApplication.java | 12 +++ .../securities/config/KafkaConfig.java | 19 ++++ .../config/SecuritiesServiceImdgConfig.java | 48 ++++++++++ .../element/SecuritiesServiceSettings.java | 31 ++++++ .../cud/MoneyMarketSecurityService.java | 27 ++++++ .../src/main/resources/application.properties | 12 +++ .../src/main/resources/logback.xml | 38 ++++++++ .../config/KafkaConsumerFactory.java | 23 +++++ .../config/KafkaProducerFactory.java | 4 +- .../config/element/KafkaConsumerSettings.java | 50 ++++++++++ ...ttings.java => KafkaProducerSettings.java} | 2 +- .../platform/messaging/domain/Consts.java | 6 ++ .../messaging/service/QueueConsumer.java | 96 +++++++++++++++++++ 18 files changed, 455 insertions(+), 8 deletions(-) create mode 100644 clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/MessagesConfig.java create mode 100644 clearing-parent/securities-service/pom.xml create mode 100644 clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/SecuritiesServiceApplication.java create mode 100644 clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/KafkaConfig.java create mode 100644 clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/SecuritiesServiceImdgConfig.java create mode 100644 clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/element/SecuritiesServiceSettings.java create mode 100644 clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java create mode 100644 clearing-parent/securities-service/src/main/resources/application.properties create mode 100644 clearing-parent/securities-service/src/main/resources/logback.xml create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaConsumerFactory.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/element/KafkaConsumerSettings.java rename platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/element/{KafkaSettings.java => KafkaProducerSettings.java} (97%) create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/MessagesConfig.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/MessagesConfig.java new file mode 100644 index 000000000..bed0c3aa7 --- /dev/null +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/MessagesConfig.java @@ -0,0 +1,21 @@ +package ru.spcex.clearing.backendapi.config; + +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.support.ResourceBundleMessageSource; + +import java.util.Locale; + +@Configuration +public class MessagesConfig { + @Bean + public ResourceBundleMessageSource messages() { + ResourceBundleMessageSource source = new ResourceBundleMessageSource(); + source.setBasenames("messages/response"); + source.setUseCodeAsDefaultMessage(true); + source.setDefaultEncoding("utf8"); + source.setDefaultLocale(Locale.ROOT); + + return source; + } +} diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/element/BackendApiSettings.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/element/BackendApiSettings.java index f4180e7c2..cdd40615c 100644 --- a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/element/BackendApiSettings.java +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/element/BackendApiSettings.java @@ -3,7 +3,7 @@ package ru.spcex.clearing.backendapi.config.element; import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.context.annotation.PropertySource; import org.springframework.stereotype.Component; -import ru.spcex.clearing.platform.messaging.config.element.KafkaSettings; +import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings; import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams; @Component @@ -11,7 +11,7 @@ import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams; @ConfigurationProperties("backend-api") public class BackendApiSettings { private HazelcastClientParams hazelcast; - private KafkaSettings kafka; + private KafkaProducerSettings kafka; private String exampleSetting; public HazelcastClientParams getHazelcast() { @@ -22,11 +22,11 @@ public class BackendApiSettings { this.hazelcast = hazelcast; } - public KafkaSettings getKafka() { + public KafkaProducerSettings getKafka() { return kafka; } - public void setKafka(KafkaSettings kafka) { + public void setKafka(KafkaProducerSettings kafka) { this.kafka = kafka; } diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/meta/CudMetaService.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/meta/CudMetaService.java index e0cef96dc..94ab34872 100644 --- a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/meta/CudMetaService.java +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/meta/CudMetaService.java @@ -3,6 +3,7 @@ package ru.spcex.clearing.backendapi.meta; import org.springframework.stereotype.Service; import ru.spcex.clearing.backendapi.controller.request.cud.MoneyMarketCreateAction; import ru.spcex.clearing.backendapi.domain.actions.IAction; +import ru.spcex.clearing.platform.messaging.domain.Consts; import java.util.HashMap; import java.util.Map; @@ -13,7 +14,7 @@ public class CudMetaService { public CudMetaService() { this.mapping = new HashMap<>(); - this.mapping.put("money-market-security-new", MoneyMarketCreateAction.class); + this.mapping.put(Consts.DESTINATION_MONEY_MARKET_SECURITY_NEW, MoneyMarketCreateAction.class); } public > Class byDestination(String destination) { diff --git a/clearing-parent/pom.xml b/clearing-parent/pom.xml index 795356502..3c0360f3e 100644 --- a/clearing-parent/pom.xml +++ b/clearing-parent/pom.xml @@ -23,6 +23,7 @@ db-scripts dbf-importer dbf-exporter + securities-service diff --git a/clearing-parent/securities-service/pom.xml b/clearing-parent/securities-service/pom.xml new file mode 100644 index 000000000..541d07d98 --- /dev/null +++ b/clearing-parent/securities-service/pom.xml @@ -0,0 +1,62 @@ + + + + clearing-parent + ru.spcex.clearing + SPCEX-1.0.0.0 + + 4.0.0 + + securities-service + + + 17 + 17 + + + + ru.spcex.platform + platform-messaging + + + ru.spcex.platform + platform-imdg-api-hazelcast-impl + + + org.springframework.boot + spring-boot-starter + + + + + + + src/main/resources + + application.properties + + false + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + repackage + + + + + ${project.artifactId} + + + + + + + \ No newline at end of file diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/SecuritiesServiceApplication.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/SecuritiesServiceApplication.java new file mode 100644 index 000000000..e1611efa7 --- /dev/null +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/SecuritiesServiceApplication.java @@ -0,0 +1,12 @@ +package ru.spcex.clearing.securities; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +@SpringBootApplication +public class SecuritiesServiceApplication { + public static void main(String[] args) { + SpringApplication app = new SpringApplication(SecuritiesServiceApplication.class); + app.run(args); + } +} diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/KafkaConfig.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/KafkaConfig.java new file mode 100644 index 000000000..deb72e8ca --- /dev/null +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/KafkaConfig.java @@ -0,0 +1,19 @@ +package ru.spcex.clearing.securities.config; + +import org.apache.kafka.clients.consumer.Consumer; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; +import ru.spcex.clearing.securities.config.element.SecuritiesServiceSettings; + +@Configuration +public class KafkaConfig { + + @Autowired + @Bean + public Consumer createProducer(SecuritiesServiceSettings settings) { + return KafkaConsumerFactory.consumer(settings.getKafka()); + } + +} diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/SecuritiesServiceImdgConfig.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/SecuritiesServiceImdgConfig.java new file mode 100644 index 000000000..4f6c0417c --- /dev/null +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/SecuritiesServiceImdgConfig.java @@ -0,0 +1,48 @@ +package ru.spcex.clearing.securities.config; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; +import ru.spcex.clearing.securities.config.element.SecuritiesServiceSettings; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; + +@Configuration +public class SecuritiesServiceImdgConfig { + @Bean(name = "taskExecutorHazelcastClientInitializer") + public ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer() { + return createThreadPoolTaskExecutor(1, true); + } + + @Bean(name = "taskExecutorIdGeneratorAwaiter") + public ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter() { + return createThreadPoolTaskExecutor(1, false); + } + + @Autowired + @Bean + public ImdgProvider imdgProvider( + @Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer, + @Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter, + SecuritiesServiceSettings settings + ) { + return new HazelcastService(taskExecutorHazelcastClientInitializer, + taskExecutorIdGeneratorAwaiter, + settings.getHazelcast()); + } + + + private static ThreadPoolTaskExecutor createThreadPoolTaskExecutor(int maxPoolSz, boolean waitForCompletion) { + ThreadPoolTaskExecutor pool = new ThreadPoolTaskExecutor(); + if (maxPoolSz > 2) { + pool.setKeepAliveSeconds(60); + pool.setAllowCoreThreadTimeOut(true); + } + pool.setCorePoolSize(maxPoolSz); + pool.setWaitForTasksToCompleteOnShutdown(waitForCompletion); + return pool; + } + +} diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/element/SecuritiesServiceSettings.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/element/SecuritiesServiceSettings.java new file mode 100644 index 000000000..546f25cb6 --- /dev/null +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/element/SecuritiesServiceSettings.java @@ -0,0 +1,31 @@ +package ru.spcex.clearing.securities.config.element; + +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.context.annotation.PropertySource; +import org.springframework.stereotype.Component; +import ru.spcex.clearing.platform.messaging.config.element.KafkaConsumerSettings; +import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams; + +@Component +@PropertySource("file:${spring.config.location}/application.properties") +@ConfigurationProperties("securities-service") +public class SecuritiesServiceSettings { + private HazelcastClientParams hazelcast; + private KafkaConsumerSettings kafka; + + public HazelcastClientParams getHazelcast() { + return hazelcast; + } + + public void setHazelcast(HazelcastClientParams hazelcast) { + this.hazelcast = hazelcast; + } + + public KafkaConsumerSettings getKafka() { + return kafka; + } + + public void setKafka(KafkaConsumerSettings kafka) { + this.kafka = kafka; + } +} diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java new file mode 100644 index 000000000..45544ab74 --- /dev/null +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java @@ -0,0 +1,27 @@ +package ru.spcex.clearing.securities.service.cud; + +import org.apache.kafka.clients.consumer.Consumer; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Service; +import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketCreateRequest; +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; + +@Service +public class MoneyMarketSecurityService extends QueueConsumer implements InitializingBean { + @Autowired + public MoneyMarketSecurityService(Consumer kafkaQueue) { + super(kafkaQueue); + } + + @Override + public void afterPropertiesSet() { + callback(Consts.DESTINATION_MONEY_MARKET_SECURITY_NEW, MoneyMarketCreateRequest.class) + .define(this::newMoneyMarket); + } + + private void newMoneyMarket(MoneyMarketCreateRequest req) { + + } +} diff --git a/clearing-parent/securities-service/src/main/resources/application.properties b/clearing-parent/securities-service/src/main/resources/application.properties new file mode 100644 index 000000000..efa54f24f --- /dev/null +++ b/clearing-parent/securities-service/src/main/resources/application.properties @@ -0,0 +1,12 @@ +spring.main.web-application-type=none + +securities-service.hazelcast.cluster-members=127.0.0.1 +securities-service.hazelcast.login=dev +securities-service.hazelcast.password=dev-pass + +securities-service.kafka.bootstrap-servers=localhost:9092 +securities-service.kafka.acks=all +securities-service.kafka.retries=0 +securities-service.kafka.batch-size=16384 +securities-service.kafka.linger-ms=1 +securities-service.kafka.buffer-memory=33554432 diff --git a/clearing-parent/securities-service/src/main/resources/logback.xml b/clearing-parent/securities-service/src/main/resources/logback.xml new file mode 100644 index 000000000..f23c2dc81 --- /dev/null +++ b/clearing-parent/securities-service/src/main/resources/logback.xml @@ -0,0 +1,38 @@ + + + + + + %date{HH:mm:ss.SSS} [%thread] %-5level %class{0}:%line - %message%n + utf-8 + + + + ./logs/securities-service.log + + + %d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %class{0}:%msg%n + utf8 + + + + ./logs/securities-service.%i.log + + 1 + 10 + + + 500MB + + + + + + + + + + + + + diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaConsumerFactory.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaConsumerFactory.java new file mode 100644 index 000000000..f2c8a7fa3 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaConsumerFactory.java @@ -0,0 +1,23 @@ +package ru.spcex.clearing.platform.messaging.config; + +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.consumer.KafkaConsumer; +import ru.spcex.clearing.platform.messaging.config.element.KafkaConsumerSettings; + +import java.util.Properties; + +public class KafkaConsumerFactory { + public static Consumer consumer(KafkaConsumerSettings kafkaSettings) { + Properties props = new Properties(); + props.put("bootstrap.servers", kafkaSettings.getBootstrapServers()); + if (kafkaSettings.getGroupId() != null && kafkaSettings.getGroupId().length() > 0) { + props.put("group.id", kafkaSettings.getGroupId()); + } + props.put("enable.auto.commit", kafkaSettings.getEnableAutoCommit().toString()); + props.put("session.timeout.ms", kafkaSettings.getSessionTimeoutMs().toString()); + props.put("auto.offset.reset", kafkaSettings.getAutoOffsetReset()); + props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); + props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); + return new KafkaConsumer<>(props); + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaProducerFactory.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaProducerFactory.java index 00498178e..b8558d665 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaProducerFactory.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaProducerFactory.java @@ -3,13 +3,13 @@ package ru.spcex.clearing.platform.messaging.config; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.ProducerConfig; -import ru.spcex.clearing.platform.messaging.config.element.KafkaSettings; +import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings; import ru.spcex.clearing.platform.messaging.serialization.JsonSerializer; import java.util.Properties; public class KafkaProducerFactory { - public static Producer producer(KafkaSettings kafkaSettings) { + public static Producer producer(KafkaProducerSettings kafkaSettings) { Properties kafkaProps = new Properties(); //Assign localhost id diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/element/KafkaConsumerSettings.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/element/KafkaConsumerSettings.java new file mode 100644 index 000000000..8e6a6fc73 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/element/KafkaConsumerSettings.java @@ -0,0 +1,50 @@ +package ru.spcex.clearing.platform.messaging.config.element; + +public class KafkaConsumerSettings { + private String bootstrapServers; + private String groupId; + private Boolean enableAutoCommit; + private Integer sessionTimeoutMs; + private String autoOffsetReset; + + + public String getBootstrapServers() { + return bootstrapServers; + } + + public void setBootstrapServers(String bootstrapServers) { + this.bootstrapServers = bootstrapServers; + } + + public String getGroupId() { + return groupId; + } + + public void setGroupId(String groupId) { + this.groupId = groupId; + } + + public Boolean getEnableAutoCommit() { + return enableAutoCommit; + } + + public void setEnableAutoCommit(Boolean enableAutoCommit) { + this.enableAutoCommit = enableAutoCommit; + } + + public Integer getSessionTimeoutMs() { + return sessionTimeoutMs; + } + + public void setSessionTimeoutMs(Integer sessionTimeoutMs) { + this.sessionTimeoutMs = sessionTimeoutMs; + } + + public String getAutoOffsetReset() { + return autoOffsetReset; + } + + public void setAutoOffsetReset(String autoOffsetReset) { + this.autoOffsetReset = autoOffsetReset; + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/element/KafkaSettings.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/element/KafkaProducerSettings.java similarity index 97% rename from platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/element/KafkaSettings.java rename to platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/element/KafkaProducerSettings.java index 27442b8d5..44f7dde60 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/element/KafkaSettings.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/element/KafkaProducerSettings.java @@ -1,6 +1,6 @@ package ru.spcex.clearing.platform.messaging.config.element; -public class KafkaSettings { +public class KafkaProducerSettings { private String bootstrapServers; private String acks; private Integer retries; diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java new file mode 100644 index 000000000..5d274e6a6 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java @@ -0,0 +1,6 @@ +package ru.spcex.clearing.platform.messaging.domain; + +public interface Consts { + String DESTINATION_MONEY_MARKET_SECURITY_NEW = "money-market-security-new"; + +} 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 new file mode 100644 index 000000000..084b3dfce --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java @@ -0,0 +1,96 @@ +package ru.spcex.clearing.platform.messaging.service; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.consumer.ConsumerRecords; +import org.apache.kafka.common.errors.WakeupException; + +import java.time.Duration; +import java.time.temporal.ChronoUnit; +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.atomic.AtomicBoolean; + +public class QueueConsumer implements AutoCloseable { + private final AtomicBoolean closed = new AtomicBoolean(false); + private final Consumer consumer; + private final ExecutorService executor; + private final Map> callbacks; + private final ObjectMapper json; + + public QueueConsumer(Consumer kafkaQueue) { + this.consumer = kafkaQueue; + this.callbacks = new HashMap<>(); + this.executor = Executors.newSingleThreadExecutor(); + this.json = new ObjectMapper(); + } + + protected void addCallBack(String destination, ConsumerWithClass callback) { + this.callbacks.put(destination, callback); + } + + public void init() throws Exception { + executor.submit(() -> { + try { + consumer.subscribe(callbacks.keySet()); + while (!closed.get()) { + ConsumerRecords records = consumer.poll(Duration.of(10, ChronoUnit.SECONDS)); + for (ConsumerRecord next : records) { + ConsumerWithClass callback = callbacks.get(next.topic()); + Object o = json.readValue((String) next.value(), callback.getClazz()); + callback.acceptRaw(o); + } + } + } catch (WakeupException e) { + if (!closed.get()) throw e; + } catch (JsonProcessingException e) { + throw new RuntimeException(e); + } finally { + consumer.close(); + } + }); + } + + @Override + public void close() throws Exception { + closed.set(true); + consumer.wakeup(); + } + + protected ConsumerWithClass callback(String destination, Class clazz) { + ConsumerWithClass callback = new ConsumerWithClass<>(clazz); + this.callbacks.put(destination, callback); + return callback; + } + + protected static class ConsumerWithClass { + private final Class clazz; + private java.util.function.Consumer consumer; + private ConsumerWithClass(Class clazz) { + this.clazz = clazz; + } + + + + public ConsumerWithClass define(java.util.function.Consumer consumer) { + this.consumer = consumer; + return this; + } + + public void accept(T obj) { + this.consumer.accept(obj); + } + + public void acceptRaw(Object obj) { + this.consumer.accept((T) obj); + } + + public Class getClazz() { + return clazz; + } + } +}