lim-exporter http://jira.mfd.msk:8088/browse/CLS-262 fix send notification
This commit is contained in:
parent
f52cc62e08
commit
ddda96e150
7 changed files with 52 additions and 11 deletions
|
|
@ -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();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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) {
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
}
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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());
|
||||
|
||||
|
|
|
|||
|
|
@ -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());
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue