diff --git a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/KafkaConfig.java b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/KafkaConfig.java index 45fc6642c..6721788f6 100644 --- a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/KafkaConfig.java +++ b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/KafkaConfig.java @@ -7,9 +7,17 @@ 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 org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.core.ProducerFactory; +import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.lim.exporter.config.settings.ExportLimServiceSettings; import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; +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 KafkaConfig { @@ -25,4 +33,24 @@ public class KafkaConfig { public Producer createProducer(ExportLimServiceSettings settings) { return KafkaProducerFactory.producer(settings.getKafkaProducer()); } + + @Bean + public KafkaTemplate kafkaTemplate(ProducerFactory pf) { + return new KafkaTemplate<>(pf); + } + + @Autowired + @Bean + public KafkaSender kafkaSender(KafkaTemplate kafkaTemplate, ImdgProvider imdgProvider) { + ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); + return KafkaSender + .setup() + .setKafkaTemplate(kafkaTemplate) + .idGenerator(imdgIdGenerator::nextId) + .imdgProvider(s -> { + Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); + return imdg::insert; + }) + .build(); + } } diff --git a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/AbstractExporterService.java b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/AbstractExporterService.java index e92739815..b26a02678 100644 --- a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/AbstractExporterService.java +++ b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/AbstractExporterService.java @@ -10,6 +10,8 @@ import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.lim.exporter.config.SFTPConfig; import ru.spcex.clearing.platform.messaging.domain.cud.utilities.LimExportedRequest; +import ru.spcex.clearing.platform.messaging.serialization.LogFormatter; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; import ru.spcex.platform.utils.log.ExceptionUtils; @@ -32,11 +34,12 @@ public abstract class AbstractExporterService { private final Imdg tradingClearingRegistryImdg; private final DateTimeFormatter dtFormatter = DateTimeFormatter.ofPattern("yyyyMMddHHmmss"); private final SFTPConfig.LimGateway gateway; - private final Producer producer; + private final KafkaSender kafkaSender; - protected AbstractExporterService(SFTPConfig.LimGateway gateway, Producer producer, ImdgProvider imdgProvider) { + protected AbstractExporterService(SFTPConfig.LimGateway gateway, + KafkaSender kafkaSender, ImdgProvider imdgProvider) { this.gateway = gateway; - this.producer = producer; + this.kafkaSender = kafkaSender; this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); this.tradingClearingRegistryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class); } @@ -68,9 +71,14 @@ public abstract class AbstractExporterService { } log.debug("Successfully exported {} file", fileName); + sendLimExportedNotification(fileName); + } + + private void sendLimExportedNotification(String fileName) { LimExportedRequest limExportedRequest = new LimExportedRequest(); limExportedRequest.setLimFileName(fileName); - producer.send(new ProducerRecord<>(LIM_EXPORTED, limExportedRequest)); + log.debug("Send message to kafka \"{}\": {}", LIM_EXPORTED, LogFormatter.toStringWrapper(limExportedRequest)); + kafkaSender.sendRequestToQueue(LIM_EXPORTED, limExportedRequest); } protected String prepareFileName(String target) { diff --git a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/MoneyExporterService.java b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/MoneyExporterService.java index 121b32210..5ab00cc30 100644 --- a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/MoneyExporterService.java +++ b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/MoneyExporterService.java @@ -1,11 +1,11 @@ package ru.spcex.clearing.lim.exporter.services; -import org.apache.kafka.clients.producer.Producer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.stereotype.Service; import ru.clearing.classes.statics.data.registry.Registry; import ru.spcex.clearing.lim.exporter.config.SFTPConfig; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.platform.imdg.api.ImdgProvider; import java.math.BigDecimal; @@ -21,9 +21,9 @@ public class MoneyExporterService extends AbstractExporterService { private final Logger log = LoggerFactory.getLogger(getClass()); public MoneyExporterService(SFTPConfig.LimGateway gateway, - Producer producer, + KafkaSender kafkaSender, ImdgProvider imdgProvider) { - super(gateway, producer, imdgProvider); + super(gateway, kafkaSender, imdgProvider); } @Override diff --git a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/SecurityExporterService.java b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/SecurityExporterService.java index 63d439dda..ae079aacc 100644 --- a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/SecurityExporterService.java +++ b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/SecurityExporterService.java @@ -6,6 +6,7 @@ import org.slf4j.LoggerFactory; import org.springframework.stereotype.Service; import ru.clearing.classes.statics.data.registry.Registry; import ru.spcex.clearing.lim.exporter.config.SFTPConfig; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.platform.imdg.api.ImdgProvider; import java.time.LocalDate; @@ -19,9 +20,9 @@ public class SecurityExporterService extends AbstractExporterService { private final Logger log = LoggerFactory.getLogger(getClass()); public SecurityExporterService(SFTPConfig.LimGateway gateway, - Producer producer, + KafkaSender kafkaSender, ImdgProvider imdgProvider) { - super(gateway, producer, imdgProvider); + super(gateway, kafkaSender, imdgProvider); } diff --git a/clearing-parent/lim-exporter/src/test/java/ru/spcex/clearing/lim/exporter/AbstractServiceTest.java b/clearing-parent/lim-exporter/src/test/java/ru/spcex/clearing/lim/exporter/AbstractServiceTest.java index 103207f84..34bb55987 100644 --- a/clearing-parent/lim-exporter/src/test/java/ru/spcex/clearing/lim/exporter/AbstractServiceTest.java +++ b/clearing-parent/lim-exporter/src/test/java/ru/spcex/clearing/lim/exporter/AbstractServiceTest.java @@ -9,6 +9,7 @@ import org.springframework.test.context.junit.jupiter.SpringExtension; import ru.clearing.classes.statics.data.registry.Registry; import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.test.config.ImdgTestConfig; import ru.spcex.clearing.test.config.KafkaTestConfig; import ru.spcex.platform.imdg.api.Imdg; @@ -37,6 +38,9 @@ public abstract class AbstractServiceTest { @Qualifier("mockProducer") protected Producer mockProducer; + @Autowired + protected KafkaSender kafkaSender; + @Autowired @Qualifier("hazelcastServiceTest") protected ImdgProvider imdgProvider; diff --git a/clearing-parent/lim-exporter/src/test/java/ru/spcex/clearing/lim/exporter/services/MoneyExporterServiceTest.java b/clearing-parent/lim-exporter/src/test/java/ru/spcex/clearing/lim/exporter/services/MoneyExporterServiceTest.java index 8163fd83e..10741693c 100644 --- a/clearing-parent/lim-exporter/src/test/java/ru/spcex/clearing/lim/exporter/services/MoneyExporterServiceTest.java +++ b/clearing-parent/lim-exporter/src/test/java/ru/spcex/clearing/lim/exporter/services/MoneyExporterServiceTest.java @@ -36,7 +36,7 @@ class MoneyExporterServiceTest extends AbstractServiceTest { registryImdg.insert(registryA); registryImdg.insert(registryD); - MoneyExporterService moneyExporterService = new MoneyExporterService(null, mockProducer, imdgProvider); + MoneyExporterService moneyExporterService = new MoneyExporterService(null, kafkaSender, imdgProvider); Collection limFileRows = moneyExporterService.getLimFileRows(); assertEquals(0, limFileRows.size()); diff --git a/clearing-parent/lim-exporter/src/test/java/ru/spcex/clearing/lim/exporter/services/SecurityExporterServiceTest.java b/clearing-parent/lim-exporter/src/test/java/ru/spcex/clearing/lim/exporter/services/SecurityExporterServiceTest.java index cf7099a4d..72813baf7 100644 --- a/clearing-parent/lim-exporter/src/test/java/ru/spcex/clearing/lim/exporter/services/SecurityExporterServiceTest.java +++ b/clearing-parent/lim-exporter/src/test/java/ru/spcex/clearing/lim/exporter/services/SecurityExporterServiceTest.java @@ -36,7 +36,7 @@ class SecurityExporterServiceTest extends AbstractServiceTest { registryImdg.insert(registry1); registryImdg.insert(registry2); - SecurityExporterService securityExporterService = new SecurityExporterService(null, mockProducer, imdgProvider); + SecurityExporterService securityExporterService = new SecurityExporterService(null, kafkaSender, imdgProvider); Collection limFileRows = securityExporterService.getLimFileRows(); assertEquals(0, limFileRows.size());