swt-exporter http://jira.mfd.msk:8088/browse/CLS-317 начало
This commit is contained in:
parent
53e10b7527
commit
d3a9c46234
18 changed files with 877 additions and 0 deletions
|
|
@ -38,6 +38,7 @@
|
|||
<module>cleaning-builders</module>
|
||||
<module>trade-importer</module>
|
||||
<module>lim-exporter</module>
|
||||
<module>swt-exporter</module>
|
||||
</modules>
|
||||
|
||||
<properties>
|
||||
|
|
|
|||
112
clearing-parent/swt-exporter/pom.xml
Normal file
112
clearing-parent/swt-exporter/pom.xml
Normal file
|
|
@ -0,0 +1,112 @@
|
|||
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
|
||||
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/maven-v4_0_0.xsd">
|
||||
<modelVersion>4.0.0</modelVersion>
|
||||
<artifactId>swt-exporter</artifactId>
|
||||
<name>Swt exporter</name>
|
||||
<version>SPCEX-1.0.0.0</version>
|
||||
<packaging>jar</packaging>
|
||||
|
||||
<parent>
|
||||
<groupId>ru.spcex.clearing</groupId>
|
||||
<artifactId>clearing-parent</artifactId>
|
||||
<version>SPCEX-1.0.0.0</version>
|
||||
</parent>
|
||||
|
||||
<properties>
|
||||
<maven.compiler.source>17</maven.compiler.source>
|
||||
<maven.compiler.target>17</maven.compiler.target>
|
||||
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
|
||||
</properties>
|
||||
|
||||
<dependencies>
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-starter</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>com.fasterxml.jackson.core</groupId>
|
||||
<artifactId>jackson-databind</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.integration</groupId>
|
||||
<artifactId>spring-integration-sftp</artifactId>
|
||||
</dependency>
|
||||
|
||||
<dependency>
|
||||
<groupId>ru.spcex.clearing</groupId>
|
||||
<artifactId>classes</artifactId>
|
||||
<version>SPCEX-1.0.0.0</version>
|
||||
<scope>compile</scope>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>ru.spcex.platform</groupId>
|
||||
<artifactId>platform-messaging</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>ru.spcex.platform</groupId>
|
||||
<artifactId>platform-imdg-api-hazelcast-impl</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>ru.spcex.platform</groupId>
|
||||
<artifactId>platform-enum</artifactId>
|
||||
</dependency>
|
||||
|
||||
<!-- TEST -->
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-starter-test</artifactId>
|
||||
<scope>test</scope>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>ru.spcex.clearing</groupId>
|
||||
<artifactId>test-clearing</artifactId>
|
||||
<scope>test</scope>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
|
||||
<build>
|
||||
<resources>
|
||||
<resource>
|
||||
<directory>src/main/resources</directory>
|
||||
<excludes>
|
||||
<exclude>application.properties</exclude>
|
||||
</excludes>
|
||||
<filtering>false</filtering>
|
||||
</resource>
|
||||
</resources>
|
||||
<plugins>
|
||||
<plugin>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-maven-plugin</artifactId>
|
||||
<executions>
|
||||
<execution>
|
||||
<goals>
|
||||
<goal>repackage</goal>
|
||||
</goals>
|
||||
</execution>
|
||||
</executions>
|
||||
<configuration>
|
||||
<finalName>${project.artifactId}</finalName>
|
||||
</configuration>
|
||||
</plugin>
|
||||
<plugin>
|
||||
<groupId>org.apache.maven.plugins</groupId>
|
||||
<artifactId>maven-surefire-plugin</artifactId>
|
||||
<version>2.21.0</version>
|
||||
<dependencies>
|
||||
<dependency>
|
||||
<groupId>org.junit.platform</groupId>
|
||||
<artifactId>junit-platform-surefire-provider</artifactId>
|
||||
<version>1.2.0-M1</version>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.junit.jupiter</groupId>
|
||||
<artifactId>junit-jupiter-engine</artifactId>
|
||||
<version>5.2.0-M1</version>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
</plugin>
|
||||
</plugins>
|
||||
</build>
|
||||
|
||||
</project>
|
||||
|
|
@ -0,0 +1,18 @@
|
|||
package ru.spcex.clearing.swt.exporter;
|
||||
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.boot.autoconfigure.SpringBootApplication;
|
||||
import org.springframework.boot.builder.SpringApplicationBuilder;
|
||||
|
||||
@SpringBootApplication
|
||||
public class SwtExportApplication {
|
||||
public static void main(String[] args) {
|
||||
try {
|
||||
SpringApplicationBuilder builder = new SpringApplicationBuilder(SwtExportApplication.class);
|
||||
builder.run(args);
|
||||
} catch (Throwable e) {
|
||||
LoggerFactory.getLogger(SwtExportApplication.class).error("Swt-exporter start failed: {} -> {}", e.getClass().getSimpleName(), e.getMessage());
|
||||
System.exit(-1);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,44 @@
|
|||
package ru.spcex.clearing.swt.exporter.config;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Qualifier;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
|
||||
import ru.spcex.clearing.swt.exporter.config.settings.ExportSwtServiceSettings;
|
||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
|
||||
|
||||
@Configuration
|
||||
public class ImdgConfig {
|
||||
private static ThreadPoolTaskExecutor createThreadPoolTaskExecutor(int maxPoolSz, boolean waitForCompletion) {
|
||||
ThreadPoolTaskExecutor pool = new ThreadPoolTaskExecutor();
|
||||
if (maxPoolSz > 2) {
|
||||
pool.setKeepAliveSeconds(60);
|
||||
pool.setAllowCoreThreadTimeOut(true);
|
||||
}
|
||||
pool.setCorePoolSize(maxPoolSz);
|
||||
pool.setWaitForTasksToCompleteOnShutdown(waitForCompletion);
|
||||
return pool;
|
||||
}
|
||||
|
||||
@Bean(name = "taskExecutorHazelcastClientInitializer")
|
||||
public ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer() {
|
||||
return createThreadPoolTaskExecutor(1, true);
|
||||
}
|
||||
|
||||
@Bean(name = "taskExecutorIdGeneratorAwaiter")
|
||||
public ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter() {
|
||||
return createThreadPoolTaskExecutor(1, false);
|
||||
}
|
||||
|
||||
@Bean("imdgProvider")
|
||||
public ImdgProvider imdgProvider(
|
||||
@Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer,
|
||||
@Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter,
|
||||
ExportSwtServiceSettings settings) {
|
||||
return new HazelcastService(taskExecutorHazelcastClientInitializer,
|
||||
taskExecutorIdGeneratorAwaiter,
|
||||
settings.getHazelcast());
|
||||
}
|
||||
|
||||
}
|
||||
|
|
@ -0,0 +1,56 @@
|
|||
package ru.spcex.clearing.swt.exporter.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 org.springframework.kafka.core.KafkaTemplate;
|
||||
import org.springframework.kafka.core.ProducerFactory;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.swt.exporter.config.settings.ExportSwtServiceSettings;
|
||||
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 {
|
||||
@Autowired
|
||||
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
|
||||
@Bean
|
||||
public Consumer<String, Object> createConsumer(ExportSwtServiceSettings settings) {
|
||||
return KafkaConsumerFactory.consumer(settings.getKafkaConsumer());
|
||||
}
|
||||
|
||||
@Autowired
|
||||
@Bean
|
||||
public Producer<String, Object> createProducer(ExportSwtServiceSettings 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();
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,11 @@
|
|||
package ru.spcex.clearing.swt.exporter.config;
|
||||
|
||||
import org.springframework.boot.context.properties.EnableConfigurationProperties;
|
||||
import org.springframework.context.annotation.ComponentScan;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
|
||||
@Configuration
|
||||
@EnableConfigurationProperties
|
||||
@ComponentScan(basePackages = {"ru.spcex.clearing.swt.exporter"})
|
||||
public class SwtExporterConfig {
|
||||
}
|
||||
|
|
@ -0,0 +1,32 @@
|
|||
package ru.spcex.clearing.swt.exporter.config.settings;
|
||||
|
||||
public class Common {
|
||||
|
||||
private String encoding;
|
||||
private int insertBatchSize;
|
||||
private int threadsCount;
|
||||
|
||||
public String getEncoding() {
|
||||
return encoding;
|
||||
}
|
||||
|
||||
public void setEncoding(String encoding) {
|
||||
this.encoding = encoding;
|
||||
}
|
||||
|
||||
public int getInsertBatchSize() {
|
||||
return insertBatchSize;
|
||||
}
|
||||
|
||||
public void setInsertBatchSize(int insertBatchSize) {
|
||||
this.insertBatchSize = insertBatchSize;
|
||||
}
|
||||
|
||||
public int getThreadsCount() {
|
||||
return threadsCount;
|
||||
}
|
||||
|
||||
public void setThreadsCount(int threadsCount) {
|
||||
this.threadsCount = threadsCount;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,59 @@
|
|||
package ru.spcex.clearing.swt.exporter.config.settings;
|
||||
|
||||
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;
|
||||
|
||||
@Component
|
||||
@PropertySource("file:${spring.config.location}/application.properties")
|
||||
@ConfigurationProperties("export-swt-service")
|
||||
public class ExportSwtServiceSettings {
|
||||
private HazelcastClientParams hazelcast;
|
||||
private KafkaConsumerSettings kafkaConsumer;
|
||||
private KafkaProducerSettings kafkaProducer;
|
||||
private Store sFTPStore;
|
||||
private Common common;
|
||||
|
||||
public HazelcastClientParams getHazelcast() {
|
||||
return hazelcast;
|
||||
}
|
||||
|
||||
public void setHazelcast(HazelcastClientParams hazelcast) {
|
||||
this.hazelcast = hazelcast;
|
||||
}
|
||||
|
||||
public KafkaConsumerSettings getKafkaConsumer() {
|
||||
return kafkaConsumer;
|
||||
}
|
||||
|
||||
public void setKafkaConsumer(KafkaConsumerSettings kafkaConsumer) {
|
||||
this.kafkaConsumer = kafkaConsumer;
|
||||
}
|
||||
|
||||
public Store getStore() {
|
||||
return sFTPStore;
|
||||
}
|
||||
|
||||
public void setStore(Store Store) {
|
||||
this.sFTPStore = Store;
|
||||
}
|
||||
|
||||
public KafkaProducerSettings getKafkaProducer() {
|
||||
return kafkaProducer;
|
||||
}
|
||||
|
||||
public void setKafkaProducer(KafkaProducerSettings kafkaProducer) {
|
||||
this.kafkaProducer = kafkaProducer;
|
||||
}
|
||||
|
||||
public Common getCommon() {
|
||||
return common;
|
||||
}
|
||||
|
||||
public void setCommon(Common common) {
|
||||
this.common = common;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,50 @@
|
|||
package ru.spcex.clearing.swt.exporter.config.settings;
|
||||
|
||||
public class Store {
|
||||
|
||||
private String outDir;
|
||||
private String user;
|
||||
private String password;
|
||||
private String serverIp;
|
||||
private int serverPort;
|
||||
|
||||
public String getUser() {
|
||||
return user;
|
||||
}
|
||||
|
||||
public void setUser(String user) {
|
||||
this.user = user;
|
||||
}
|
||||
|
||||
public String getPassword() {
|
||||
return password;
|
||||
}
|
||||
|
||||
public void setPassword(String password) {
|
||||
this.password = password;
|
||||
}
|
||||
|
||||
public String getServerIp() {
|
||||
return serverIp;
|
||||
}
|
||||
|
||||
public void setServerIp(String serverIp) {
|
||||
this.serverIp = serverIp;
|
||||
}
|
||||
|
||||
public int getServerPort() {
|
||||
return serverPort;
|
||||
}
|
||||
|
||||
public void setServerPort(int serverPort) {
|
||||
this.serverPort = serverPort;
|
||||
}
|
||||
|
||||
public String getOutDir() {
|
||||
return outDir;
|
||||
}
|
||||
|
||||
public void setOutDir(String outDir) {
|
||||
this.outDir = outDir;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,78 @@
|
|||
package ru.spcex.clearing.swt.exporter.services;
|
||||
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
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.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 java.io.IOException;
|
||||
import java.nio.file.Files;
|
||||
import java.time.LocalDateTime;
|
||||
import java.time.format.DateTimeFormatter;
|
||||
import java.util.Collection;
|
||||
import java.util.List;
|
||||
|
||||
import static ru.spcex.clearing.platform.messaging.domain.Consts.LIM_EXPORTED;
|
||||
|
||||
public abstract class AbstractExporterService {
|
||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||
protected final Imdg<Registry> registryImdg;
|
||||
private final List<String> validStatus = List.of("ACTV", "ROPN");
|
||||
private final Imdg<TradingClearingRegistry> tradingClearingRegistryImdg;
|
||||
private final DateTimeFormatter dtFormatter = DateTimeFormatter.ofPattern("yyyyMMddHHmmss");
|
||||
private final KafkaSender kafkaSender;
|
||||
protected final FileStorage fileStorage;
|
||||
|
||||
protected AbstractExporterService(FileStorage fileStorage,
|
||||
KafkaSender kafkaSender, ImdgProvider imdgProvider) {
|
||||
this.fileStorage = fileStorage;
|
||||
this.kafkaSender = kafkaSender;
|
||||
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
|
||||
this.tradingClearingRegistryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
|
||||
}
|
||||
|
||||
public abstract Collection<String> getLimFileRows();
|
||||
|
||||
public abstract String getTargetFileName();
|
||||
|
||||
public void process() {
|
||||
String fileName = getTargetFileName();
|
||||
log.debug("Start export {} Lim file", fileName);
|
||||
|
||||
try {
|
||||
// todo возможная оптимизация: посмотреть размеры файлов, возможно обойтись без временного файла
|
||||
//Files.write(limFilePath, getLimFileRows());
|
||||
fileStorage.saveFile(fileName, null);
|
||||
} catch (IOException e) {
|
||||
log.error("Failed export {} file", fileName);
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
log.debug("Successfully exported {} file", fileName);
|
||||
|
||||
sendSwtxportedNotification(fileName);
|
||||
}
|
||||
|
||||
void sendSwtxportedNotification(String fileName) {
|
||||
LimExportedRequest limExportedRequest = new LimExportedRequest();
|
||||
limExportedRequest.setLimFileName(fileName);
|
||||
log.debug("Send message to kafka \"{}\": {}", LIM_EXPORTED, LogFormatter.toStringWrapper(limExportedRequest));
|
||||
kafkaSender.sendRequestToQueue(LIM_EXPORTED, limExportedRequest);//todo rewrite!!!
|
||||
}
|
||||
|
||||
protected String prepareFileName(Long counter, String target) {
|
||||
String dt = dtFormatter.format(LocalDateTime.now());
|
||||
String type="09";
|
||||
String section="U";
|
||||
String counterS = counter==null?"":"_"+counter;
|
||||
String partyCode="";
|
||||
String result="KS_RDC_DF-%s_%s_PRC%s%s%s.swt".formatted(type, section,dt,counterS,partyCode);
|
||||
return result;
|
||||
}
|
||||
|
||||
}
|
||||
|
|
@ -0,0 +1,35 @@
|
|||
package ru.spcex.clearing.swt.exporter.services;
|
||||
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.stereotype.Service;
|
||||
import ru.spcex.clearing.swt.exporter.config.settings.ExportSwtServiceSettings;
|
||||
|
||||
import java.io.File;
|
||||
import java.io.IOException;
|
||||
|
||||
@Service
|
||||
public class FileStorage {
|
||||
protected final Logger log= LoggerFactory.getLogger(getClass());
|
||||
protected File outPath;
|
||||
|
||||
@Autowired
|
||||
public FileStorage(ExportSwtServiceSettings config) {
|
||||
if (config.getStore().getOutDir()==null || config.getStore().getOutDir().isBlank()) {
|
||||
throw new IllegalArgumentException("Out directory settings is empty.");
|
||||
}
|
||||
this.outPath = new File(config.getStore().getOutDir()); //todo ...
|
||||
if (!outPath.isDirectory()) {
|
||||
log.info("Path not exist. mkdir \"{}\"", outPath.getAbsolutePath());
|
||||
if (!outPath.mkdir()) {
|
||||
log.error("Can not make output directory \"{}\"", outPath);
|
||||
}
|
||||
}
|
||||
log.info("Output directory \"{}\"", outPath);
|
||||
}
|
||||
|
||||
public void saveFile(String fileName, byte[] data) throws IOException {
|
||||
//todo ...
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,38 @@
|
|||
package ru.spcex.clearing.swt.exporter.services;
|
||||
|
||||
import org.apache.kafka.clients.consumer.Consumer;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.beans.factory.InitializingBean;
|
||||
import org.springframework.stereotype.Service;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
|
||||
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||
import ru.spcex.clearing.swt.exporter.services.exportimpl.MoneyExporterService;
|
||||
import ru.spcex.platform.enumeration.Task;
|
||||
|
||||
@Service
|
||||
public class LauncherCommandReceiver extends QueueConsumer implements InitializingBean {
|
||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||
private final MoneyExporterService moneyExporterService;
|
||||
// private final SecurityExporterService securityExporterService;
|
||||
|
||||
public LauncherCommandReceiver(Consumer<String, Object> kafkaQueue,
|
||||
MoneyExporterService moneyExporterService
|
||||
// , SecurityExporterService securityExporterService
|
||||
) {
|
||||
super(kafkaQueue);
|
||||
this.moneyExporterService = moneyExporterService;
|
||||
// this.securityExporterService = securityExporterService;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void afterPropertiesSet() {
|
||||
callback(LauncherCommandRequest.class)
|
||||
.setConsumer(action -> moneyExporterService.process())
|
||||
.forDestination(Task.unloadingSession_LIMM.topic(), callbacks::put); // LIMM
|
||||
// callback(LauncherCommandRequest.class)
|
||||
// .setConsumer(action -> securityExporterService.process())
|
||||
// .forDestination(Task.unloadingSession_LIMS.topic(), callbacks::put); // LIMS
|
||||
init();
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,103 @@
|
|||
package ru.spcex.clearing.swt.exporter.services.exportimpl;
|
||||
|
||||
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.platform.messaging.service.sender.KafkaSender;
|
||||
import ru.spcex.clearing.swt.exporter.services.AbstractExporterService;
|
||||
import ru.spcex.clearing.swt.exporter.services.FileStorage;
|
||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||
|
||||
import java.math.BigDecimal;
|
||||
import java.time.LocalDate;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collection;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
@Service
|
||||
public class MoneyExporterService extends AbstractExporterService {
|
||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||
|
||||
public MoneyExporterService(FileStorage fileStorage,
|
||||
KafkaSender kafkaSender,
|
||||
ImdgProvider imdgProvider) {
|
||||
super(fileStorage, kafkaSender, imdgProvider);
|
||||
}
|
||||
|
||||
@Override
|
||||
public String getTargetFileName() {
|
||||
return prepareFileName(null, "money");
|
||||
}
|
||||
|
||||
@Override
|
||||
public Collection<String> getLimFileRows() {
|
||||
log.debug("Started loading and formation of money file lines");
|
||||
LocalDate currentDate = LocalDate.now();
|
||||
List<String> swtFileRows = new ArrayList<>();
|
||||
Collection<Registry> registriesA = registryImdg.getCollectionObjectsByFieldValues(Map.of(
|
||||
"registryDesignation", "A",
|
||||
"registryInstrumentType", "M",
|
||||
"registryUnit", "F"
|
||||
));
|
||||
Collection<Registry> registriesD = registryImdg.getCollectionObjectsByFieldValues(Map.of(
|
||||
"registryDesignation", "D",
|
||||
"registryInstrumentType", "M",
|
||||
"registryUnit", "T"
|
||||
));
|
||||
Map<String, List<Registry>> byTcrA = registriesA.stream()
|
||||
.collect(Collectors.groupingBy(Registry::getTradingClearingRegistry));
|
||||
Map<String, List<Registry>> byTcrD = registriesD.stream()
|
||||
.collect(Collectors.groupingBy(Registry::getTradingClearingRegistry));
|
||||
|
||||
for (Map.Entry<String, List<Registry>> entryA : byTcrA.entrySet()) {
|
||||
List<Registry> registriesListA = entryA.getValue();
|
||||
List<Registry> registriesListB = byTcrD.get(entryA.getKey());
|
||||
for (Registry registryA : registriesListA) {
|
||||
// todo if (checkNotBlocked(registryA)) {
|
||||
// Registry registryD = findRegistryBySecurityId(registryA.getSecurityId(), registriesListB);
|
||||
// limFileRows.add(getRow(registryA, registryD));
|
||||
// }
|
||||
}
|
||||
}
|
||||
log.debug("Successfully completed the formation of rows: {} for export money", swtFileRows.size());
|
||||
return swtFileRows;
|
||||
}
|
||||
|
||||
public String getRow(Registry registryA, Registry registryD) {
|
||||
StringBuilder row = new StringBuilder();
|
||||
|
||||
row.append("MONEY: FIRM_ID = ");
|
||||
row.append(registryA.getTradingCode());
|
||||
|
||||
row.append("; TAG = SPVB");
|
||||
|
||||
row.append("; CURR_CODE = ");
|
||||
row.append(registryA.getSecuritySymbol());
|
||||
|
||||
row.append("; CLIENT_CODE = ");
|
||||
row.append(registryA.getTradingClearingRegistry());
|
||||
|
||||
row.append("; OPEN_BALANCE = ");
|
||||
BigDecimal balance = registryA.getBalance() != null ?
|
||||
registryD != null && registryD.getBalance() != null ? registryA.getBalance().subtract(registryD.getBalance()) : registryA.getBalance() :
|
||||
BigDecimal.ZERO;
|
||||
row.append(balance);
|
||||
|
||||
row.append("; OPEN_LIMIT = 0.00");
|
||||
|
||||
row.append("; LIMIT_KIND = 0;");
|
||||
return row.toString();
|
||||
}
|
||||
|
||||
private Registry findRegistryBySecurityId(Long securityId, List<Registry> registries) {
|
||||
for (Registry registry : registries) {
|
||||
if (securityId.equals(registry.getSecurityId())) {
|
||||
return registry;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,25 @@
|
|||
spring.main.web-application-type=none
|
||||
|
||||
export-swt-service.hazelcast.cluster-members=10.200.200.181:5701
|
||||
export-swt-service.hazelcast.login=dev
|
||||
export-swt-service.hazelcast.password=dev-pass
|
||||
|
||||
export-swt-service.common.encoding=cp866
|
||||
export-swt-service.common.threads-count=10
|
||||
|
||||
export-swt-service.out-dir=DocOut
|
||||
|
||||
export-swt-service.kafka-consumer.bootstrap-servers=localhost:9092
|
||||
export-swt-service.kafka-consumer.group-id=dev-group-balance-service
|
||||
export-swt-service.kafka-consumer.enable-auto-commit=false
|
||||
export-swt-service.kafka-consumer.session-timeout-ms=30000
|
||||
export-swt-service.kafka-consumer.auto-offset-reset=latest
|
||||
export-swt-service.kafka-consumer.linger-ms=1
|
||||
export-swt-service.kafka-consumer.buffer-memory=33554432
|
||||
|
||||
export-swt-service.kafka-producer.bootstrap-servers=localhost:9092
|
||||
export-swt-service.kafka-producer.acks=all
|
||||
export-swt-service.kafka-producer.retries=0
|
||||
export-swt-service.kafka-producer.batch-size=16384
|
||||
export-swt-service.kafka-producer.linger-ms=1
|
||||
export-swt-service.kafka-producer.buffer-memory=33554432
|
||||
37
clearing-parent/swt-exporter/src/main/resources/logback.xml
Normal file
37
clearing-parent/swt-exporter/src/main/resources/logback.xml
Normal file
|
|
@ -0,0 +1,37 @@
|
|||
<configuration>
|
||||
|
||||
<appender name="CONSOLE" class="ch.qos.logback.core.ConsoleAppender">
|
||||
<encoder>
|
||||
<charset>UTF-8</charset>
|
||||
<pattern>%date{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n</pattern>
|
||||
</encoder>
|
||||
</appender>
|
||||
|
||||
<appender name="FILE" class="ch.qos.logback.core.rolling.RollingFileAppender">
|
||||
<file>./logs/swt-exporter.log</file>
|
||||
<encoder>
|
||||
<charset>UTF-8</charset>
|
||||
<pattern>%d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n</pattern>
|
||||
</encoder>
|
||||
<rollingPolicy class="ch.qos.logback.core.rolling.FixedWindowRollingPolicy">
|
||||
<fileNamePattern>
|
||||
../logs/swt-exporter.%i.log
|
||||
</fileNamePattern>
|
||||
<minIndex>1</minIndex>
|
||||
<maxIndex>10</maxIndex>
|
||||
</rollingPolicy>
|
||||
<triggeringPolicy class="ch.qos.logback.core.rolling.SizeBasedTriggeringPolicy">
|
||||
<maxFileSize>500MB</maxFileSize>
|
||||
</triggeringPolicy>
|
||||
</appender>
|
||||
|
||||
<root level="warn">
|
||||
<appender-ref ref="CONSOLE"/>
|
||||
<appender-ref ref="FILE"/>
|
||||
</root>
|
||||
|
||||
<logger name="ru.spcex" level="debug" additivity="false">
|
||||
<appender-ref ref="FILE"/>
|
||||
<appender-ref ref="CONSOLE"/>
|
||||
</logger>
|
||||
</configuration>
|
||||
|
|
@ -0,0 +1,53 @@
|
|||
package ru.spcex.clearing.swt.exporter;
|
||||
|
||||
import org.apache.kafka.clients.producer.Producer;
|
||||
import org.junit.jupiter.api.extension.ExtendWith;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.beans.factory.annotation.Qualifier;
|
||||
import org.springframework.test.context.ContextConfiguration;
|
||||
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;
|
||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||
|
||||
import java.time.LocalDate;
|
||||
|
||||
import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID;
|
||||
import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId;
|
||||
|
||||
@ExtendWith(SpringExtension.class)
|
||||
@ContextConfiguration(classes = {
|
||||
ImdgTestConfig.class,
|
||||
KafkaTestConfig.class})
|
||||
public abstract class AbstractServiceTest {
|
||||
protected static final long id = currentID.getAndIncrement();
|
||||
protected Imdg<Registry> registryImdg;
|
||||
protected Imdg<TradingClearingRegistry> tradingClearingRegistryImdg;
|
||||
protected LocalDate currentDate = LocalDate.now();
|
||||
protected String tcrA = "1324A234";
|
||||
protected String tcrD = "124324A234";
|
||||
protected Long securityIdFirst = 12L;
|
||||
protected Long securityIdSecond = 23L;
|
||||
|
||||
@Autowired
|
||||
@Qualifier("mockProducer")
|
||||
protected Producer<String, Object> mockProducer;
|
||||
|
||||
@Autowired
|
||||
protected KafkaSender kafkaSender;
|
||||
|
||||
@Autowired
|
||||
@Qualifier("hazelcastServiceTest")
|
||||
protected ImdgProvider imdgProvider;
|
||||
|
||||
protected void init() {
|
||||
waitAvailableImdgProviderAndAddAdminWithDefaultId();
|
||||
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
|
||||
this.tradingClearingRegistryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,23 @@
|
|||
package ru.spcex.clearing.swt.exporter.services;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import ru.spcex.clearing.swt.exporter.AbstractServiceTest;
|
||||
import ru.spcex.clearing.swt.exporter.services.exportimpl.MoneyExporterService;
|
||||
|
||||
import javax.annotation.PostConstruct;
|
||||
|
||||
class AbstractExporterServiceTest extends AbstractServiceTest {
|
||||
@PostConstruct
|
||||
public void init() {
|
||||
super.init();
|
||||
}
|
||||
|
||||
@Test
|
||||
void sendSwtExportedNotification() {
|
||||
AbstractExporterService moneyExporterService = new MoneyExporterService(null, kafkaSender, imdgProvider);
|
||||
|
||||
String fileName = "KS_RDC_DF-14_fund_202305241832.swt";
|
||||
moneyExporterService.sendSwtxportedNotification(fileName);
|
||||
//TestUtils.waitingSendAndCheckRecord(null, mockProducer);
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,102 @@
|
|||
package ru.spcex.clearing.swt.exporter.services;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.springframework.util.StringUtils;
|
||||
import ru.clearing.classes.statics.data.registry.Registry;
|
||||
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
|
||||
import ru.spcex.clearing.swt.exporter.AbstractServiceTest;
|
||||
import ru.spcex.clearing.swt.exporter.services.exportimpl.MoneyExporterService;
|
||||
|
||||
import javax.annotation.PostConstruct;
|
||||
import java.math.BigDecimal;
|
||||
import java.util.Collection;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
import static ru.spcex.clearing.test.TestUtils.clearAllInImdg;
|
||||
|
||||
class MoneyExporterServiceTest extends AbstractServiceTest {
|
||||
@PostConstruct
|
||||
public void init() {
|
||||
super.init();
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link MoneyExporterService#getSwtFileRows()}<br>
|
||||
* Тест проверяет создание строк документа lim.<br>
|
||||
*/
|
||||
@Test // todo rewrite
|
||||
void getSwtFileRows() {
|
||||
clearAllInImdg(tradingClearingRegistryImdg);
|
||||
Registry registryA = getRegistryA(tcrA, securityIdFirst);
|
||||
Registry registryD = getRegistryD(tcrA, securityIdFirst);
|
||||
registryImdg.insert(registryA);
|
||||
registryImdg.insert(registryD);
|
||||
registryA = getRegistryA(tcrD, securityIdSecond);
|
||||
registryD = getRegistryD(tcrD, securityIdSecond);
|
||||
registryImdg.insert(registryA);
|
||||
registryImdg.insert(registryD);
|
||||
|
||||
MoneyExporterService moneyExporterService = new MoneyExporterService(null, kafkaSender, imdgProvider);
|
||||
Collection<String> limFileRows = moneyExporterService.getLimFileRows();
|
||||
assertEquals(0, limFileRows.size());
|
||||
|
||||
TradingClearingRegistry tradingClearingRegistry = new TradingClearingRegistry();
|
||||
tradingClearingRegistry.setId(securityIdFirst);
|
||||
tradingClearingRegistry.setStatus("ACTV");
|
||||
tradingClearingRegistryImdg.insert(tradingClearingRegistry);
|
||||
tradingClearingRegistry.setId(securityIdSecond);
|
||||
tradingClearingRegistryImdg.insert(tradingClearingRegistry);
|
||||
|
||||
limFileRows = moneyExporterService.getLimFileRows();
|
||||
assertEquals(2, limFileRows.size());
|
||||
assertTrue(limFileRows.contains(moneyExporterService.getRow(registryA, registryD)));
|
||||
}
|
||||
|
||||
private Registry getRegistryA(String tradingClearingRegistry, Long securityId) {
|
||||
Registry registry = new Registry();
|
||||
registry.setTradingCode("1A12323");
|
||||
registry.setBalance(new BigDecimal("10.00"));
|
||||
registry.setTradingClearingRegistry(tradingClearingRegistry);
|
||||
registry.setRegistryDesignation("A");
|
||||
registry.setRegistryInstrumentType("M");
|
||||
registry.setRegistryUnit("F");
|
||||
registry.setClearingCode(clearingCode(registry));
|
||||
registry.setClearingDate(currentDate);
|
||||
registry.setSecuritySymbol("RUB");
|
||||
registry.setSecurityId(securityId);
|
||||
registry.setTradingClearingRegistryId(securityId);
|
||||
return registry;
|
||||
}
|
||||
|
||||
private Registry getRegistryD(String tradingClearingRegistry, Long securityId) {
|
||||
Registry registry = new Registry();
|
||||
registry.setTradingCode("1A12323");
|
||||
registry.setBalance(new BigDecimal("5.00"));
|
||||
registry.setTradingClearingRegistry(tradingClearingRegistry);
|
||||
registry.setRegistryDesignation("D");
|
||||
registry.setRegistryInstrumentType("M");
|
||||
registry.setRegistryUnit("T");
|
||||
registry.setClearingCode(clearingCode(registry));
|
||||
registry.setClearingDate(currentDate);
|
||||
registry.setSecuritySymbol("RUB");
|
||||
registry.setSecurityId(securityId);
|
||||
registry.setTradingClearingRegistryId(securityId);
|
||||
return registry;
|
||||
}
|
||||
|
||||
|
||||
// see clearing-service ReistryUtil:
|
||||
|
||||
public static String clearingCode(Registry ofRegistry) {
|
||||
return clearingCode(ofRegistry.getRegistryDesignation(), ofRegistry.getRegistryInstrumentType(), ofRegistry.getRegistryCapacity(), ofRegistry.getRegistryUnit());
|
||||
}
|
||||
|
||||
public static String clearingCode(String registryDesignation, String registryInstrumentType, String registryCapacity, String registryUnit) {
|
||||
if (StringUtils.isEmpty(registryDesignation)) registryDesignation = "-";
|
||||
if (StringUtils.isEmpty(registryInstrumentType)) registryInstrumentType = "-";
|
||||
if (StringUtils.isEmpty(registryCapacity)) registryCapacity = "-";
|
||||
if (StringUtils.isEmpty(registryUnit)) registryUnit = "-";
|
||||
return registryDesignation + registryInstrumentType + registryCapacity + registryUnit;
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue