From ae3b24c8a1d074f45e84eec7525aae8d72688abd Mon Sep 17 00:00:00 2001 From: AKurakin Date: Mon, 26 Dec 2022 19:30:17 +0300 Subject: [PATCH] =?UTF-8?q?reports-service=20=D0=B4=D0=BE=D0=B1=D0=B0?= =?UTF-8?q?=D0=B2=D0=BB=D0=B5=D0=BD=D0=B8=D0=B5=20kafka,=20=D0=BD=D0=B0?= =?UTF-8?q?=D1=87=D0=B0=D0=BB=D0=BE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- clearing-parent/reports-service/pom.xml | 4 ++ .../clearing/reports/config/KafkaConfig.java | 28 +++++++++++ .../element/ReportsServiceSettings.java | 22 ++++++--- .../reports/services/QCommandExecutor.java | 49 +++++++++++++++++++ .../src/main/resources/application.properties | 17 +++++++ 5 files changed, 114 insertions(+), 6 deletions(-) create mode 100644 clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/config/KafkaConfig.java create mode 100644 clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/QCommandExecutor.java diff --git a/clearing-parent/reports-service/pom.xml b/clearing-parent/reports-service/pom.xml index 54846a228..457291ada 100644 --- a/clearing-parent/reports-service/pom.xml +++ b/clearing-parent/reports-service/pom.xml @@ -62,6 +62,10 @@ ru.spcex.platform platform-imdg-api-hazelcast-impl + + ru.spcex.platform + platform-messaging + org.junit.jupiter diff --git a/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/config/KafkaConfig.java b/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/config/KafkaConfig.java new file mode 100644 index 000000000..cd3b73821 --- /dev/null +++ b/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/config/KafkaConfig.java @@ -0,0 +1,28 @@ +package ru.spcex.clearing.reports.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; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Scope; +import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; +import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; +import ru.spcex.clearing.reports.config.element.ReportsServiceSettings; + +@Configuration +public class KafkaConfig { + @Autowired + @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) + @Bean + public Consumer createConsumer(ReportsServiceSettings settings) { + return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); + } + +// @Autowired +// @Bean +// public Producer createProducer(ReportsServiceSettings settings) { +// return KafkaProducerFactory.producer(settings.getKafkaProducer()); +// } +} diff --git a/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/config/element/ReportsServiceSettings.java b/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/config/element/ReportsServiceSettings.java index a79eaa702..d3aa03d8a 100644 --- a/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/config/element/ReportsServiceSettings.java +++ b/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/config/element/ReportsServiceSettings.java @@ -3,6 +3,8 @@ package ru.spcex.clearing.reports.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.clearing.platform.messaging.config.element.KafkaProducerSettings; import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams; /** @@ -15,7 +17,8 @@ import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams; @ConfigurationProperties("reports-service") public class ReportsServiceSettings { private HazelcastClientParams hazelcast; - private String reportModuleOut; + private KafkaConsumerSettings kafkaConsumer; + private KafkaProducerSettings kafkaProducer; public HazelcastClientParams getHazelcast() { return hazelcast; @@ -25,12 +28,19 @@ public class ReportsServiceSettings { this.hazelcast = hazelcast; } - - public String getReportModuleOut() { - return reportModuleOut; + public KafkaConsumerSettings getKafkaConsumer() { + return kafkaConsumer; } - public void setReportModuleOut(String reportModuleOut) { - this.reportModuleOut = reportModuleOut; + public void setKafkaConsumer(KafkaConsumerSettings kafkaConsumer) { + this.kafkaConsumer = kafkaConsumer; + } + + public KafkaProducerSettings getKafkaProducer() { + return kafkaProducer; + } + + public void setKafkaProducer(KafkaProducerSettings kafkaProducer) { + this.kafkaProducer = kafkaProducer; } } diff --git a/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/QCommandExecutor.java b/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/QCommandExecutor.java new file mode 100644 index 000000000..a093f6eec --- /dev/null +++ b/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/QCommandExecutor.java @@ -0,0 +1,49 @@ +package ru.spcex.clearing.reports.services; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.InitializingBean; +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; + +import org.apache.kafka.clients.consumer.Consumer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.domain.Consts; +//import ru.spcex.clearing.platform.messaging.domain.cud.balance.ExportToFileRequest; todo use class? +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.enumeration.Task; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgId; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.util.List; + +public class QCommandExecutor extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); +// private final ImdgId idGenerator; +// private final KafkaSender kafkaReqProducer; + + + public QCommandExecutor(Consumer kafkaQueue, ImdgProvider imdgProvider, + List listOfCollectors) { + super(kafkaQueue); +// todo constructor + } + + @Override + public void afterPropertiesSet() { + log.debug("Init queue listener {}", getClass().getSimpleName()); +// callback(Object.class) +// .setConsumer(this::newReportDo) +// .forDestination(Task.getVerification!!!.topic(), callbacks::put); // todo здесь надо указать команду из ТЗ или согласовать с фронтэндом +// init(); + } + + private void newReportDo(BaseRequest userRequest) { + log.debug("newReportDo request received"); + // todo logic + log.debug("successfully processed"); + } +} diff --git a/clearing-parent/reports-service/src/main/resources/application.properties b/clearing-parent/reports-service/src/main/resources/application.properties index 52d18bea7..c7765a7fe 100644 --- a/clearing-parent/reports-service/src/main/resources/application.properties +++ b/clearing-parent/reports-service/src/main/resources/application.properties @@ -2,7 +2,24 @@ server.port=8080 server.servlet.context-path=/reports_service spring.main.web-application-type=servlet + +reports-service.kafka-consumer.bootstrap-servers=localhost:9092 +reports-service.kafka-consumer.group-id=dev-group-balance-service +reports-service.kafka-consumer.enable-auto-commit=false +reports-service.kafka-consumer.session-timeout-ms=30000 +reports-service.kafka-consumer.auto-offset-reset=latest +reports-service.kafka-consumer.linger-ms=1 +reports-service.kafka-consumer.buffer-memory=33554432 + +#reports-service.kafka-producer.bootstrap-servers=localhost:9092 +#reports-service.kafka-producer.acks=all +#reports-service.kafka-producer.retries=0 +#reports-service.kafka-producer.batch-size=16384 +#reports-service.kafka-producer.linger-ms=1 +#reports-service.kafka-producer.buffer-memory=33554432 + reports-service.hazelcast.cluster-members=127.0.0.1:5701,10.200.200.181:5701 reports-service.hazelcast.login=dev reports-service.hazelcast.password=dev-pass + reports-service.ReportModuleOut=./report_out