From 7f542d6736be7f9f30e95e566dae021b3b409c63 Mon Sep 17 00:00:00 2001 From: ialbert Date: Fri, 26 May 2023 16:56:08 +0300 Subject: [PATCH] kafkaSender for backend-api --- .../backendapi/config/KafkaConfig.java | 42 +++++++++++++++++++ 1 file changed, 42 insertions(+) diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/KafkaConfig.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/KafkaConfig.java index 9046c9bc8..409aee180 100644 --- a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/KafkaConfig.java +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/KafkaConfig.java @@ -1,5 +1,8 @@ package ru.spcex.clearing.backendapi.config; +import org.apache.kafka.clients.admin.AdminClient; +import org.apache.kafka.clients.admin.AdminClientConfig; +import org.apache.kafka.clients.admin.NewTopic; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.MockConsumer; import org.apache.kafka.clients.consumer.OffsetResetStrategy; @@ -9,9 +12,23 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Profile; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.core.ProducerFactory; import ru.spcex.clearing.backendapi.config.element.BackendApiSettings; import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; +import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings; +import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.imdg.api.ImdgId; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import javax.annotation.PostConstruct; +import java.lang.reflect.Field; +import java.util.ArrayList; +import java.util.List; +import java.util.Properties; +import java.util.concurrent.ExecutionException; @Configuration public class KafkaConfig { @@ -23,6 +40,31 @@ public class KafkaConfig { return KafkaProducerFactory.producer(settings.getKafkaProducer()); } + + + @Bean + public ProducerFactory pf(BackendApiSettings settings) { + KafkaProducerSettings kafkaSettings = settings.getKafkaProducer(); + return KafkaProducerFactory.producerFactory(kafkaSettings); + } + + @Bean + public KafkaTemplate kafkaTemplateS(ProducerFactory pf) { + return new KafkaTemplate<>(pf); + } + + @Profile("!kafkaDisabled") + @Autowired + @Bean + public KafkaSender kafkaSender(KafkaTemplate kafkaTemplate, ImdgProvider imdgProvider) { + ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); + return KafkaSender.setup() + .saveRequestInfo(false) + .setKafkaTemplate(kafkaTemplate) + .idGenerator(imdgIdGenerator::nextId) + .build(); + } + @Profile("kafkaDisabled") @Bean public MockProducer createMockProducer() {