diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/KafkaConfig.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/KafkaConfig.java index 1cfeb4a56..d1c48b183 100644 --- a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/KafkaConfig.java +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/KafkaConfig.java @@ -1,6 +1,7 @@ package ru.spcex.clearing.balance.config; import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.config.ConfigurableBeanFactory; import org.springframework.context.annotation.Bean; @@ -8,13 +9,20 @@ import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Scope; import ru.spcex.clearing.balance.config.element.BalanceServiceSettings; import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; +import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; @Configuration public class KafkaConfig { @Autowired @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) @Bean - public Consumer createProducer(BalanceServiceSettings settings) { - return KafkaConsumerFactory.consumer(settings.getKafka()); + public Consumer createConsumer(BalanceServiceSettings settings) { + return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); + } + + @Autowired + @Bean + public Producer createProducer(BalanceServiceSettings settings) { + return KafkaProducerFactory.producer(settings.getKafkaProducer()); } } diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/KafkaSenderConfig.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/KafkaSenderConfig.java new file mode 100644 index 000000000..fa621992c --- /dev/null +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/KafkaSenderConfig.java @@ -0,0 +1,30 @@ +package ru.spcex.clearing.balance.config; + +import org.apache.kafka.clients.producer.Producer; +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.service.RequestInfo; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgId; +import ru.spcex.platform.imdg.api.ImdgProvider; + +@Configuration +public class KafkaSenderConfig { + + @Autowired + @Bean + public KafkaSender kafkaSender(Producer kafkaProducer, ImdgProvider imdgProvider) { + ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); + return KafkaSender + .setup() + .producer(kafkaProducer) + .idGenerator(imdgIdGenerator::nextId) + .imdgProvider(s -> { + Imdg imdg = imdgProvider.getImdg(s, RequestInfo.class); + return imdg::insert; + }) + .build(); + } +} diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/element/BalanceServiceSettings.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/element/BalanceServiceSettings.java index dbdb8c45a..f4136355f 100644 --- a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/element/BalanceServiceSettings.java +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/element/BalanceServiceSettings.java @@ -4,6 +4,7 @@ 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.clearing.platform.messaging.config.element.KafkaProducerSettings; import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams; @Component @@ -11,7 +12,8 @@ import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams; @ConfigurationProperties("balance-service") public class BalanceServiceSettings { private HazelcastClientParams hazelcast; - private KafkaConsumerSettings kafka; + private KafkaConsumerSettings kafkaConsumer; + private KafkaProducerSettings kafkaProducer; public HazelcastClientParams getHazelcast() { return hazelcast; @@ -21,11 +23,19 @@ public class BalanceServiceSettings { this.hazelcast = hazelcast; } - public KafkaConsumerSettings getKafka() { - return kafka; + public KafkaConsumerSettings getKafkaConsumer() { + return kafkaConsumer; } - public void setKafka(KafkaConsumerSettings kafka) { - this.kafka = kafka; + public void setKafkaConsumer(KafkaConsumerSettings kafkaConsumer) { + this.kafkaConsumer = kafkaConsumer; + } + + public KafkaProducerSettings getKafkaProducer() { + return kafkaProducer; + } + + public void setKafkaProducer(KafkaProducerSettings kafkaProducer) { + this.kafkaProducer = kafkaProducer; } } diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaImdgInsert.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaImdgInsert.java new file mode 100644 index 000000000..241501e0f --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaImdgInsert.java @@ -0,0 +1,7 @@ +package ru.spcex.clearing.platform.messaging.service.sender; + +import ru.spcex.clearing.platform.messaging.service.RequestInfo; + +public interface KafkaImdgInsert { + void insert(RequestInfo requestInfo); +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSender.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSender.java new file mode 100644 index 000000000..8bb67b5ec --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSender.java @@ -0,0 +1,76 @@ +package ru.spcex.clearing.platform.messaging.service.sender; + +import org.apache.kafka.clients.producer.Producer; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.clients.producer.RecordMetadata; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import ru.spcex.clearing.platform.messaging.domain.ActionType; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.service.RequestInfo; +import ru.spcex.platform.classes.base.SpcexObjectBase; +import ru.spcex.platform.utils.log.ExceptionUtils; + +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Future; +import java.util.function.Function; +import java.util.function.Supplier; + +public class KafkaSender { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Map allImdgMaps; + private Producer kafka; + private Supplier idGenerator; + private Function imdgGenerator; + + KafkaSender() { + this.allImdgMaps = new ConcurrentHashMap<>(); + } + + public static KafkaSenderBuilderImpl setup() { + return new KafkaSenderBuilderImpl(); + } + + + public Long sendRequestToQueue(String destination, Object requestPayload) { + BaseRequest request = new BaseRequest<>(); + request.setId(idGenerator.get()); + request.setActionType(ActionType.SYSTEM); + request.setRequestPayload(requestPayload); + //сохраняет данные о запросе в хранилище + saveRequestToStorage(destination, request); + Future send = kafka.send(new ProducerRecord<>(destination, request)); + try { + send.get(); + } catch (InterruptedException | ExecutionException e) { + log.error(ExceptionUtils.getStackTrace(e)); + return null; + } + return request.getId(); + } + + private void saveRequestToStorage(String destination, BaseRequest request) { + KafkaImdgInsert imdgInsert = getImdg(destination); + RequestInfo requestInfo = RequestInfo.create(request.getId()); + imdgInsert.insert(requestInfo); + } + + @SuppressWarnings("unchecked") + private KafkaImdgInsert getImdg(String mapName) { + return allImdgMaps.computeIfAbsent(mapName, (mapName1) -> imdgGenerator.apply(mapName)); + } + + void setProducer(Producer kafka) { + this.kafka = kafka; + } + + void setIdGenerator(Supplier idGenerator) { + this.idGenerator = idGenerator; + } + + void setImdgProvider(Function imdgGenerator) { + this.imdgGenerator = imdgGenerator; + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSenderBuilder.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSenderBuilder.java new file mode 100644 index 000000000..f07f23463 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSenderBuilder.java @@ -0,0 +1,16 @@ +package ru.spcex.clearing.platform.messaging.service.sender; + +import org.apache.kafka.clients.producer.Producer; + +import java.util.function.Function; +import java.util.function.Supplier; + +public interface KafkaSenderBuilder { + KafkaSenderBuilder producer(Producer kafka); + + KafkaSenderBuilder idGenerator(Supplier idGenerator); + + KafkaSenderBuilder imdgProvider(Function imdgGenerator); + + KafkaSender build(); +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSenderBuilderImpl.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSenderBuilderImpl.java new file mode 100644 index 000000000..d20b848ec --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSenderBuilderImpl.java @@ -0,0 +1,38 @@ +package ru.spcex.clearing.platform.messaging.service.sender; + +import org.apache.kafka.clients.producer.Producer; + +import java.util.function.Function; +import java.util.function.Supplier; + +public class KafkaSenderBuilderImpl implements KafkaSenderBuilder{ + + private final KafkaSender kafkaSender; + + public KafkaSenderBuilderImpl() { + this.kafkaSender = new KafkaSender(); + } + + @Override + public KafkaSenderBuilder producer(Producer kafka) { + kafkaSender.setProducer(kafka); + return this; + } + + @Override + public KafkaSenderBuilder idGenerator(Supplier idGenerator) { + kafkaSender.setIdGenerator(idGenerator); + return this; + } + + @Override + public KafkaSenderBuilder imdgProvider(Function imdgGenerator) { + kafkaSender.setImdgProvider(imdgGenerator); + return this; + } + + @Override + public KafkaSender build() { + return kafkaSender; + } +}