diff --git a/clearing-parent/lim-exporter/pom.xml b/clearing-parent/lim-exporter/pom.xml new file mode 100644 index 000000000..1a9f7225b --- /dev/null +++ b/clearing-parent/lim-exporter/pom.xml @@ -0,0 +1,112 @@ + + 4.0.0 + lim-exporter + Lim exporter + SPCEX-1.0.0.0 + jar + + + ru.spcex.clearing + clearing-parent + SPCEX-1.0.0.0 + + + + 17 + 17 + UTF-8 + + + + + org.springframework.boot + spring-boot-starter + + + com.fasterxml.jackson.core + jackson-databind + + + org.springframework.integration + spring-integration-sftp + + + + ru.spcex.clearing + classes + SPCEX-1.0.0.0 + compile + + + ru.spcex.platform + platform-messaging + + + ru.spcex.platform + platform-imdg-api-hazelcast-impl + + + ru.spcex.platform + platform-enum + + + + + org.springframework.boot + spring-boot-starter-test + test + + + ru.spcex.clearing + test-clearing + test + + + + + + + src/main/resources + + application.properties + + false + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + repackage + + + + + ${project.artifactId} + + + + org.apache.maven.plugins + maven-surefire-plugin + 2.21.0 + + + org.junit.platform + junit-platform-surefire-provider + 1.2.0-M1 + + + org.junit.jupiter + junit-jupiter-engine + 5.2.0-M1 + + + + + + + diff --git a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/LimExportApplication.java b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/LimExportApplication.java new file mode 100644 index 000000000..f2d29d17a --- /dev/null +++ b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/LimExportApplication.java @@ -0,0 +1,18 @@ +package ru.spcex.clearing.lim.exporter; + +import org.slf4j.LoggerFactory; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.boot.builder.SpringApplicationBuilder; + +@SpringBootApplication +public class LimExportApplication { + public static void main(String[] args) { + try { + SpringApplicationBuilder builder = new SpringApplicationBuilder(LimExportApplication.class); + builder.run(args); + } catch (Throwable e) { + LoggerFactory.getLogger(LimExportApplication.class).error("DBF-Loader start failed: {} -> {}", e.getClass().getSimpleName(), e.getMessage()); + System.exit(-1); + } + } +} diff --git a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/ImdgConfig.java b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/ImdgConfig.java new file mode 100644 index 000000000..34ed1dcd9 --- /dev/null +++ b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/ImdgConfig.java @@ -0,0 +1,46 @@ +package ru.spcex.clearing.lim.exporter.config; + +import org.springframework.beans.factory.annotation.Autowired; +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.lim.exporter.config.settings.ExportLimServiceSettings; +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); + } + + @Autowired + @Bean + public ImdgProvider imdgProvider( + @Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer, + @Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter, + ExportLimServiceSettings settings) { + return new HazelcastService(taskExecutorHazelcastClientInitializer, + taskExecutorIdGeneratorAwaiter, + settings.getHazelcast()); + } + +} 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 new file mode 100644 index 000000000..45fc6642c --- /dev/null +++ b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/KafkaConfig.java @@ -0,0 +1,28 @@ +package ru.spcex.clearing.lim.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 ru.spcex.clearing.lim.exporter.config.settings.ExportLimServiceSettings; +import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; +import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; + +@Configuration +public class KafkaConfig { + @Autowired + @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) + @Bean + public Consumer createConsumer(ExportLimServiceSettings settings) { + return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); + } + + @Autowired + @Bean + public Producer createProducer(ExportLimServiceSettings settings) { + return KafkaProducerFactory.producer(settings.getKafkaProducer()); + } +} diff --git a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/LimExporterConfig.java b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/LimExporterConfig.java new file mode 100644 index 000000000..b81b8820f --- /dev/null +++ b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/LimExporterConfig.java @@ -0,0 +1,63 @@ +package ru.spcex.clearing.lim.exporter.config; + +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.ComponentScan; +import org.springframework.context.annotation.Configuration; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; +import ru.spcex.clearing.lim.exporter.config.settings.ExportLimServiceSettings; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams; +import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; + +@Configuration +@EnableConfigurationProperties +@ComponentScan(basePackages = {"ru.spcex.clearing.lim.exporter"}) +public class LimExporterConfig { + + @Bean("taskExecutorHazelcastClientInitializer") + public ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer() { + return createThreadPoolTaskExecutor(1, true); + } + + @Bean("taskExecutorIdGeneratorAwaiter") + public ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter() { + return createThreadPoolTaskExecutor(1, false); + } + + @Bean("imdgProvider") + public ImdgProvider imdgProvider(@Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer, + @Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter, + ExportLimServiceSettings settings) { + HazelcastClientParams params = new HazelcastClientParams(); + params.setClusterMembers(settings.getHazelcast().getClusterMembers()); + params.setLogin(settings.getHazelcast().getLogin()); + params.setPassword(settings.getHazelcast().getPassword()); + return new HazelcastService(taskExecutorHazelcastClientInitializer, taskExecutorIdGeneratorAwaiter, params); + } + + @Bean("executor") + public ThreadPoolTaskExecutor executor(ExportLimServiceSettings settings) { + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setMaxPoolSize(settings.getCommon().getThreadsCount()); + executor.setCorePoolSize(settings.getCommon().getThreadsCount()); + executor.setThreadNamePrefix("dbf-exporter"); + executor.setWaitForTasksToCompleteOnShutdown(true); + executor.setAwaitTerminationSeconds(300); + executor.initialize(); + return executor; + } + + 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; + } + +} diff --git a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/SFTPConfig.java b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/SFTPConfig.java new file mode 100644 index 000000000..c6f9c1670 --- /dev/null +++ b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/SFTPConfig.java @@ -0,0 +1,54 @@ +package ru.spcex.clearing.lim.exporter.config; + +import com.jcraft.jsch.ChannelSftp; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.expression.common.LiteralExpression; +import org.springframework.integration.annotation.Gateway; +import org.springframework.integration.annotation.MessagingGateway; +import org.springframework.integration.annotation.ServiceActivator; +import org.springframework.integration.file.remote.session.CachingSessionFactory; +import org.springframework.integration.file.remote.session.SessionFactory; +import org.springframework.integration.sftp.outbound.SftpMessageHandler; +import org.springframework.integration.sftp.session.DefaultSftpSessionFactory; +import org.springframework.messaging.MessageHandler; +import ru.spcex.clearing.lim.exporter.config.settings.ExportLimServiceSettings; + +import java.io.File; + +@Configuration +public class SFTPConfig { + + @Bean + public SessionFactory sftpSessionFactory(ExportLimServiceSettings settings) { + DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory(true); + factory.setHost(settings.getStore().getServerIp()); + factory.setPort(settings.getStore().getServerPort()); + factory.setUser(settings.getStore().getUser()); + factory.setPassword(settings.getStore().getPassword()); + factory.setAllowUnknownKeys(true); + return new CachingSessionFactory<>(factory); + } + + @Bean + @ServiceActivator(inputChannel = "toSftpChannel") + public MessageHandler handler(SessionFactory sessionFactory, ExportLimServiceSettings settings) { + SftpMessageHandler handler = new SftpMessageHandler(sessionFactory); + handler.setRemoteDirectoryExpression(new LiteralExpression(settings.getStore().getOutDir())); + handler.setAutoCreateDirectory(true); + handler.setFileNameGenerator(message -> { + if (message.getPayload() instanceof File) { + return ((File) message.getPayload()).getName(); + }else { + throw new IllegalArgumentException("File expected as payload."); + } + }); + return handler; + } + + @MessagingGateway + public interface LimGateway { + @Gateway(requestChannel = "toSftpChannel") + void sendToSftp(File file); + } +} diff --git a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/settings/Common.java b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/settings/Common.java new file mode 100644 index 000000000..28dab5aa9 --- /dev/null +++ b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/settings/Common.java @@ -0,0 +1,32 @@ +package ru.spcex.clearing.lim.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; + } +} diff --git a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/settings/ExportLimServiceSettings.java b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/settings/ExportLimServiceSettings.java new file mode 100644 index 000000000..9bb6de24a --- /dev/null +++ b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/settings/ExportLimServiceSettings.java @@ -0,0 +1,59 @@ +package ru.spcex.clearing.lim.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-lim-service") +public class ExportLimServiceSettings { + 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; + } +} diff --git a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/settings/Store.java b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/settings/Store.java new file mode 100644 index 000000000..ec4c6012f --- /dev/null +++ b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/config/settings/Store.java @@ -0,0 +1,50 @@ +package ru.spcex.clearing.lim.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; + } +} 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 new file mode 100644 index 000000000..61cb29394 --- /dev/null +++ b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/AbstractExporterService.java @@ -0,0 +1,56 @@ +package ru.spcex.clearing.lim.exporter.services; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import ru.spcex.clearing.lim.exporter.config.SFTPConfig; + +import java.io.File; +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.time.LocalDateTime; +import java.time.format.DateTimeFormatter; +import java.util.Collection; + +public abstract class AbstractExporterService { + private final Logger log = LoggerFactory.getLogger(getClass()); + private static final DateTimeFormatter dtFormatter = DateTimeFormatter.ofPattern("yyyyMMddHHmmss"); + private final SFTPConfig.LimGateway gateway; + + protected AbstractExporterService(SFTPConfig.LimGateway gateway) { + this.gateway = gateway; + } + + public abstract Collection getLimFileRows(); + public abstract String getTargetFileName(); + + public void process() { + log.debug("Start export {} Lim file", getTargetFileName()); + String fileName = getTargetFileName(); + File limFile = new File(fileName); + + Path limFilePath = limFile.toPath(); + try { + Files.write(limFilePath, getLimFileRows()); + gateway.sendToSftp(limFile); + + log.debug("Successfully exported {} file", fileName); + Files.deleteIfExists(limFilePath); + } catch (IOException e) { + log.error("Failed export {} file", fileName); + throw new RuntimeException(e); + } + } + + protected String prepareFileName(String target){ + StringBuilder result = new StringBuilder(); + String dt = dtFormatter.format(LocalDateTime.now()); + + result.append("limits_"); + result.append(target); + result.append('_'); + result.append(dt); + result.append(".lim"); + return result.toString(); + } +} diff --git a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/LauncherCommandReceiver.java b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/LauncherCommandReceiver.java new file mode 100644 index 000000000..660a4c287 --- /dev/null +++ b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/LauncherCommandReceiver.java @@ -0,0 +1,36 @@ +package ru.spcex.clearing.lim.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.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 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(); + } +} 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 new file mode 100644 index 000000000..c66f61648 --- /dev/null +++ b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/MoneyExporterService.java @@ -0,0 +1,104 @@ +package ru.spcex.clearing.lim.exporter.services; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.stereotype.Service; +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.lim.exporter.config.SFTPConfig; +import ru.spcex.platform.imdg.api.Imdg; +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()); + private final Imdg registryImdg; + private final List validStatus = List.of("ACTV", "ROPN"); + private final Imdg tradingClearingRegistryImdg; + + public MoneyExporterService(SFTPConfig.LimGateway gateway, ImdgProvider imdgProvider) { + super(gateway); + this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); + this.tradingClearingRegistryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class); + } + + @Override + public String getTargetFileName() { + return prepareFileName("money"); + } + + @Override + public Collection getLimFileRows() { + log.debug("Started loading and formation of money file lines"); + LocalDate currentDate = LocalDate.now(); + List limFileRows = new ArrayList<>(); + Collection registriesA = registryImdg.getCollectionObjectsByFieldValues( + Map.of("registryDesignation", "A", + "registryInstrumentType", "M", + "registryCode", "F", + "clearingDate", currentDate)); + Collection registriesD = registryImdg.getCollectionObjectsByFieldValues( + Map.of("registryDesignation", "D", + "registryInstrumentType", "M", + "registryCode", "T", + "clearingDate", currentDate)); + Map> byTcrA = registriesA.stream() + .collect(Collectors.groupingBy(Registry::getTradingClearingRegistry)); + Map> byTcrD = registriesD.stream() + .collect(Collectors.groupingBy(Registry::getTradingClearingRegistry)); + + for(Map.Entry> entryA : byTcrA.entrySet()){ + List registriesListA = entryA.getValue(); + List registriesListB = byTcrD.get(entryA.getKey()); + for (Registry registry : registriesListA){ + if(checkNotBloked(registry)){ + //из ТЗ пока не понятно как соотнести Registry из registriesListA и registriesListB по валюте(RUB, USD) и вычислить остаток на денежном счете +// limFileRows.add(getRow(...)); + } + } + } + log.debug("Successfully completed the formation of rows: {} for export money", limFileRows.size()); + return limFileRows; + } + + private boolean checkNotBloked(Registry registry){ + TradingClearingRegistry tradingClearingRegistry = tradingClearingRegistryImdg.getSingleObjectByID(registry.getTradingClearingRegistryId()); + return validStatus.contains(tradingClearingRegistry.getStatus()); + } + + //поменяется при обновлении ТЗ + private String getRow(Registry registryA, Registry registryB){ + StringBuilder row = new StringBuilder(); + + row.append("MONEY: FIRM_ID = "); + row.append(registryA.getTradingCode()); + + row.append("; TAG = SPVB"); + + row.append("; CURR_CODE = "); +// MoneyMarketSecurity mms = moneyMarketSecurityMap.getSingleObjectByFieldValues(Map.of("securityId", )); + + row.append("; CLIENT_CODE = "); + row.append(registryA.getTradingClearingRegistry()); + + row.append("; OPEN_BALANCE = "); + BigDecimal balance = registryA.getBalance() != null ? + registryB.getBalance() != null ? registryA.getBalance().subtract(registryB.getBalance()) : registryA.getBalance() : + BigDecimal.ZERO; + row.append(balance); + + row.append("; OPEN_LIMIT = 0.00"); + + row.append("; LIMIT_KIND = 0;"); + return row.toString(); + } +} 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 new file mode 100644 index 000000000..206bbfb42 --- /dev/null +++ b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/SecurityExporterService.java @@ -0,0 +1,81 @@ +package ru.spcex.clearing.lim.exporter.services; + +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.imdg.IMDGDistributedNames; +import ru.spcex.clearing.lim.exporter.config.SFTPConfig; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.time.LocalDate; +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; +import java.util.Map; + +@Service +public class SecurityExporterService extends AbstractExporterService { + private final Logger log = LoggerFactory.getLogger(getClass()); + + private final Imdg registryImdg; + + public SecurityExporterService(SFTPConfig.LimGateway gateway, ImdgProvider imdgProvider) { + super(gateway); + this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); + } + + + @Override + public String getTargetFileName() { + return prepareFileName("security"); + } + + @Override + public Collection getLimFileRows(){ + log.debug("Started loading and formation of DEPO file lines"); + LocalDate currentDate = LocalDate.now(); + List limFileRows = new ArrayList<>(); + Collection registries = registryImdg.getCollectionObjectsByFieldValues( + Map.of("registryDesignation", "A", + "registryInstrumentType", "S", + "registryCode", "T", + "clearingDate", currentDate)); + for (Registry registry : registries){ + limFileRows.add(getRow(registry)); + } + log.debug("Successfully completed the formation of rows: {} for export DEPO", limFileRows.size()); + return limFileRows; + } + + private String getRow(Registry registry){ + StringBuilder row = new StringBuilder(); + + row.append("DEPO: FIRM_ID = "); + row.append(registry.getTradingCode()); + + row.append("; SECCODE = "); + row.append(getSecurityShortName(registry)); + + row.append("; CLIENT_CODE = "); + row.append(registry.getTradingClearingRegistry()); + + row.append("; OPEN_BALANCE = "); + row.append(registry.getBalance()); + + row.append("; OPEN_LIMIT = 0"); + + row.append("; TRDACCID = "); + row.append(registry.getTradingClearingRegistry()); + + row.append("; LIMIT_KIND = 0;"); + + return row.toString(); + } + + //скорее всего в новой редакции ТЗ не понадобится + private String getSecurityShortName(Registry registry){ + return null; + } +} diff --git a/clearing-parent/lim-exporter/src/main/resources/application.properties b/clearing-parent/lim-exporter/src/main/resources/application.properties new file mode 100644 index 000000000..114f0954c --- /dev/null +++ b/clearing-parent/lim-exporter/src/main/resources/application.properties @@ -0,0 +1,33 @@ +spring.main.web-application-type=none + +export-lim-service.hazelcast.cluster-members=10.200.200.181:5701 +export-lim-service.hazelcast.login=dev +export-lim-service.hazelcast.password=dev-pass + +export-lim-service.common.encoding=cp866 +export-lim-service.common.threads-count=10 + +export-lim-service.store.out-dir=DocOut +export-lim-service.store.user:tester +export-lim-service.store.password=password +export-lim-service.store.server-ip=10.230.238.53 +export-lim-service.store.server-port=2222 + +export-lim-service.hazelcast.cluster-members=127.0.0.1:5701 +export-lim-service.hazelcast.login=dev +export-lim-service.hazelcast.password=dev-pass + +export-lim-service.kafka-consumer.bootstrap-servers=localhost:9092 +export-lim-service.kafka-consumer.group-id=dev-group-balance-service +export-lim-service.kafka-consumer.enable-auto-commit=false +export-lim-service.kafka-consumer.session-timeout-ms=30000 +export-lim-service.kafka-consumer.auto-offset-reset=latest +export-lim-service.kafka-consumer.linger-ms=1 +export-lim-service.kafka-consumer.buffer-memory=33554432 + +export-lim-service.kafka-producer.bootstrap-servers=localhost:9092 +export-lim-service.kafka-producer.acks=all +export-lim-service.kafka-producer.retries=0 +export-lim-service.kafka-producer.batch-size=16384 +export-lim-service.kafka-producer.linger-ms=1 +export-lim-service.kafka-producer.buffer-memory=33554432 \ No newline at end of file diff --git a/clearing-parent/lim-exporter/src/main/resources/logback.xml b/clearing-parent/lim-exporter/src/main/resources/logback.xml new file mode 100644 index 000000000..d1a27ddbf --- /dev/null +++ b/clearing-parent/lim-exporter/src/main/resources/logback.xml @@ -0,0 +1,37 @@ + + + + + UTF-8 + %date{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n + + + + + ./logs/lim-exporter.log + + UTF-8 + %d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n + + + + ../logs/lim-exporter.%i.log + + 1 + 10 + + + 500MB + + + + + + + + + + + + + \ No newline at end of file diff --git a/clearing-parent/pom.xml b/clearing-parent/pom.xml index f879cdf1e..c80817211 100644 --- a/clearing-parent/pom.xml +++ b/clearing-parent/pom.xml @@ -37,6 +37,7 @@ test-clearing cleaning-builders trade-importer + lim-exporter diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java index 678b0acf5..e6b512d58 100644 --- a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java @@ -16,7 +16,8 @@ public enum Task implements IEnumKey { createOrderConfirm("CORC"), getAllBalance("GALB"), createReport_GREP("GREP"), // Создание отчёта (report-service) RPRT нескольких видов, этот GREP - unloadingSession_LIMM("LIMM"),//Выгрузка в торговую систему остатков секции МКР + unloadingSession_LIMM("LIMM"),//Выгрузка в торговую систему остатков по деньгам + unloadingSession_LIMS("LIMS"),//Выгрузка в торговую систему остатков по бумагам unloadingSession_LIMF("LIMF"),//Выгрузка в торговую систему остатков Фондовой секции liquidationSession_LIQU("LIQU"),//Ликвидационная сессия по обязательтсвам участника startSession_STRM("STRM"),//Начало торговой сессии секции МКР