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() {