Merge remote-tracking branch 'origin/dev' into dev

This commit is contained in:
ialbert 2023-05-24 15:01:57 +03:00
commit d9cc47aadf
7 changed files with 52 additions and 11 deletions

View file

@ -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<String, Object> createProducer(ExportLimServiceSettings settings) {
return KafkaProducerFactory.producer(settings.getKafkaProducer());
}
@Bean
public KafkaTemplate<String, Object> kafkaTemplate(ProducerFactory<String, Object> pf) {
return new KafkaTemplate<>(pf);
}
@Autowired
@Bean
public KafkaSender kafkaSender(KafkaTemplate<String, Object> kafkaTemplate, ImdgProvider imdgProvider) {
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
return KafkaSender
.setup()
.setKafkaTemplate(kafkaTemplate)
.idGenerator(imdgIdGenerator::nextId)
.imdgProvider(s -> {
Imdg<RequestInfo> imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);
return imdg::insert;
})
.build();
}
}

View file

@ -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<TradingClearingRegistry> tradingClearingRegistryImdg;
private final DateTimeFormatter dtFormatter = DateTimeFormatter.ofPattern("yyyyMMddHHmmss");
private final SFTPConfig.LimGateway gateway;
private final Producer<String, Object> producer;
private final KafkaSender kafkaSender;
protected AbstractExporterService(SFTPConfig.LimGateway gateway, Producer<String, Object> 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) {

View file

@ -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<String, Object> producer,
KafkaSender kafkaSender,
ImdgProvider imdgProvider) {
super(gateway, producer, imdgProvider);
super(gateway, kafkaSender, imdgProvider);
}
@Override

View file

@ -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<String, Object> producer,
KafkaSender kafkaSender,
ImdgProvider imdgProvider) {
super(gateway, producer, imdgProvider);
super(gateway, kafkaSender, imdgProvider);
}

View file

@ -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<String, Object> mockProducer;
@Autowired
protected KafkaSender kafkaSender;
@Autowired
@Qualifier("hazelcastServiceTest")
protected ImdgProvider imdgProvider;

View file

@ -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<String> limFileRows = moneyExporterService.getLimFileRows();
assertEquals(0, limFileRows.size());

View file

@ -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<String> limFileRows = securityExporterService.getLimFileRows();
assertEquals(0, limFileRows.size());