This commit is contained in:
AKurakin 2024-02-09 11:37:44 +03:00
parent f033ea8244
commit 1f2f70d301
31 changed files with 1050 additions and 1 deletions

View file

@ -219,6 +219,7 @@
<task id="27" code="EDEP" name="Время завершения возврата депозитов"/>
<task id="28" code="CHDF" name="Проверка наличия пары ДФ-01/ДФ-57 и ДФ-08/ДФ-21"/>
<task id="29" code="CCLR" name="Завершение неудачных клиринговых сессий"/>
<task id="30" code="CBRR" name="Загрузка кросс-курсов"/>
<taskStatus id="1" code="ACTV" name="Активна"/>
<taskStatus id="2" code="BLKD" name="Не активна"/>
<taskStatus id="3" code="CNCL" name="Отмена расписания"/>

View file

@ -5390,6 +5390,17 @@
"fields": []
}
,
{"method":"post",
"destination": "CBRR",
"group": "Обмен с интеграционными модулями",
"name": "Загрузка кросс-курсов",
"fields": []
}
,
{"method":"post",
"destination": "LIMM",

View file

@ -1251,6 +1251,8 @@
</post>
<post destination="GTRD" group="Обмен с Торговой системой" name="Получение сделок из Торговой системы">
</post>
<post destination="CBRR" group="Обмен с интеграционными модулями" name="Загрузка кросс-курсов">
</post>
<post destination="LIMM" group="Обмен с Торговой системой" name="Выгрузка в Торговую систему остатков по деньгам (отправка lim)">
</post>
<post destination="LIMS" group="Обмен с Торговой системой" name="Выгрузка в Торговую систему остатков по бумагам (отправка lim)">

View file

@ -5390,6 +5390,17 @@
"fields": []
}
,
{"method":"post",
"destination": "CBRR",
"group": "Обмен с интеграционными модулями",
"name": "Загрузка кросс-курсов",
"fields": []
}
,
{"method":"post",
"destination": "LIMM",

View file

@ -438,6 +438,8 @@ INSERT INTO TASK_DICTIONARY(ID, CODE, NAME) values (28, 'CHDF', 'Проверк
INSERT INTO TASK_DICTIONARY(ID, CODE, NAME) values (29, 'CCLR', 'Завершение неудачных клиринговых сессий') ON CONFLICT (ID) DO UPDATE SET CODE = EXCLUDED.CODE, NAME = EXCLUDED.NAME;
INSERT INTO TASK_DICTIONARY(ID, CODE, NAME) values (30, 'CBRR', 'Загрузка кросс-курсов') ON CONFLICT (ID) DO UPDATE SET CODE = EXCLUDED.CODE, NAME = EXCLUDED.NAME;
INSERT INTO TASK_STATUS_DICTIONARY(ID, CODE, NAME) values (1, 'ACTV', 'Активна') ON CONFLICT (ID) DO UPDATE SET CODE = EXCLUDED.CODE, NAME = EXCLUDED.NAME;
INSERT INTO TASK_STATUS_DICTIONARY(ID, CODE, NAME) values (2, 'BLKD', 'Не активна') ON CONFLICT (ID) DO UPDATE SET CODE = EXCLUDED.CODE, NAME = EXCLUDED.NAME;

View file

@ -0,0 +1,138 @@
<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>
<parent>
<groupId>ru.spcex.clearing</groupId>
<artifactId>clearing-parent</artifactId>
<version>SPCEX-1.0.0.0</version>
</parent>
<artifactId>extdb-importer</artifactId>
<name>extdb-importer</name>
<description>Extdb importer module. Аналог trade-importer.</description>
<version>SPCEX-1.0.0.0</version>
<packaging>jar</packaging>
<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>org.springframework.boot</groupId>
<artifactId>spring-boot-configuration-processor</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>com.fasterxml.jackson.core</groupId>
<artifactId>jackson-databind</artifactId>
</dependency>
<!-- JDBC -->
<dependency>
<groupId>org.springframework</groupId>
<artifactId>spring-jdbc</artifactId>
</dependency>
<dependency>
<groupId>com.microsoft.sqlserver</groupId>
<artifactId>mssql-jdbc</artifactId>
</dependency>
<dependency>
<groupId>com.mchange</groupId>
<artifactId>c3p0</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.clearing</groupId>
<artifactId>clearing-validation</artifactId>
</dependency>
<dependency>
<groupId>ru.spcex.platform</groupId>
<artifactId>platform-enum</artifactId>
</dependency>
<!-- special logging -->
<dependency>
<groupId>net.logstash.logback</groupId>
<artifactId>logstash-logback-encoder</artifactId>
<version>7.0.1</version>
</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>

View file

@ -0,0 +1,12 @@
package ru.spex.clearing.extdb.importer;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@SpringBootApplication
public class ExtDBImporterApplication {
public static void main(String[] args) {
SpringApplication app = new SpringApplication(ExtDBImporterApplication.class);
app.run(args);
}
}

View file

@ -0,0 +1,65 @@
package ru.spex.clearing.extdb.importer.config;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.jdbc.datasource.DriverManagerDataSource;
import ru.spex.clearing.extdb.importer.config.settings.DatabaseSettings;
import ru.spex.clearing.extdb.importer.config.settings.ImportExtDBServiceSettings;
import ru.spex.clearing.extdb.importer.error.ModuleInitializeException;
import javax.sql.DataSource;
import java.sql.Connection;
import static org.springframework.jdbc.datasource.DataSourceUtils.doCloseConnection;
@SuppressWarnings("UnnecessaryLocalVariable")
@Configuration
public class DbConnectionConfig {
private final Logger log = LoggerFactory.getLogger(this.getClass());
private final DatabaseSettings settings;
public DbConnectionConfig(ImportExtDBServiceSettings settings) {
this.settings = settings.getDatabase();
}
@Bean
public DataSource dataSource() {
String login = settings.getLogin();
String password = settings.getPassword();
String dbUrl = settings.getUrl();
String driver = settings.getDriver();
DriverManagerDataSource ds = new DriverManagerDataSource();
try {
ds.setDriverClassName(driver);
} catch (Exception ue) {
throw new RuntimeException("JDBC driver not loaded: " + driver, ue);
}
ds.setUrl(dbUrl);
ds.setUsername(login);
ds.setPassword(password);
String OPERATION_DATABASE_CONNECTION_CHECK = String.format("Database [%s] connection check", dbUrl);
try {
Connection conn = ds.getConnection();
doCloseConnection(conn, ds);
log.info("{}: success", OPERATION_DATABASE_CONNECTION_CHECK);
return ds;
} catch (Throwable e) {
String msg = String.format("%s: failed: %s -> %s",
OPERATION_DATABASE_CONNECTION_CHECK, e.getClass().getSimpleName(), e.getMessage());
log.error(msg);
throw new ModuleInitializeException(msg, e);
}
}
@Bean
public JdbcTemplate jdbcTemplate(DataSource dataSource) {
JdbcTemplate jdbcTemplate = new JdbcTemplate(dataSource);
return jdbcTemplate;
}
}

View file

@ -0,0 +1,15 @@
package ru.spex.clearing.extdb.importer.config;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import ru.spcex.platform.utils.enumeration.IMessageResolver;
import ru.spcex.clearing.util.services.IMDGMessageResolver;
import ru.spcex.platform.imdg.api.ImdgProvider;
@Configuration
public class ErrorResolverConfig {
@Bean
public IMessageResolver messageResolver(ImdgProvider imdgProvider) {
return new IMDGMessageResolver(imdgProvider);
}
}

View file

@ -0,0 +1,20 @@
package ru.spex.clearing.extdb.importer.config;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.context.ApplicationContext;
import org.springframework.context.annotation.ComponentScan;
import org.springframework.context.annotation.Configuration;
import ru.spex.clearing.extdb.importer.config.settings.ImportExtDBServiceSettings;
@Configuration
@EnableConfigurationProperties
@ComponentScan(basePackages = {"ru.spcex.clearing.extdb.importer"})
public class ExtDBImporterConfig {
private final ImportExtDBServiceSettings settings;
private final ApplicationContext context;
public ExtDBImporterConfig(ImportExtDBServiceSettings settings, ApplicationContext context) {
this.settings = settings;
this.context = context;
}
}

View file

@ -0,0 +1,47 @@
package ru.spex.clearing.extdb.importer.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.spex.clearing.extdb.importer.config.settings.ImportExtDBServiceSettings;
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,
ImportExtDBServiceSettings settings
) {
return new HazelcastService(taskExecutorHazelcastClientInitializer,
taskExecutorIdGeneratorAwaiter,
settings.getHazelcast());
}
}

View file

@ -0,0 +1,60 @@
package ru.spex.clearing.extdb.importer.config;
import org.apache.kafka.clients.consumer.Consumer;
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.platform.messaging.config.KafkaConsumerFactory;
import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory;
import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings;
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spex.clearing.extdb.importer.config.settings.ImportExtDBServiceSettings;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
import java.util.function.Supplier;
@Configuration
public class KafkaConfig {
@Autowired
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
@Bean
public Consumer<String, Object> createConsumer(ImportExtDBServiceSettings settings) {
return KafkaConsumerFactory.consumer(settings.getKafkaConsumer());
}
@Bean
public ProducerFactory<String, Object> pf(ImportExtDBServiceSettings settings) {
KafkaProducerSettings kafkaSettings = settings.getKafkaProducer();
return KafkaProducerFactory.producerFactory(kafkaSettings);
}
@Bean("kafkaTemplate")
public KafkaTemplate<String, Object> kafkaTemplate(ProducerFactory<String, Object> pf) {
return new KafkaTemplate<>(pf);
}
@Autowired
@Bean
public Supplier<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

@ -0,0 +1,14 @@
package ru.spex.clearing.extdb.importer.config.settings;
public class Cron {
private String checkSrcDirCron;
public String getCheckSrcDirCron() {
return checkSrcDirCron;
}
public void setCheckSrcDirCron(String checkSrcDirCron) {
this.checkSrcDirCron = checkSrcDirCron;
}
}

View file

@ -0,0 +1,40 @@
package ru.spex.clearing.extdb.importer.config.settings;
public class DatabaseSettings {
private String login;
private String password;
private String url;
private String driver;
public String getLogin() {
return login;
}
public void setLogin(String login) {
this.login = login;
}
public String getPassword() {
return password;
}
public void setPassword(String password) {
this.password = password;
}
public String getUrl() {
return url;
}
public void setUrl(String url) {
this.url = url;
}
public String getDriver() {
return driver;
}
public void setDriver(String driver) {
this.driver = driver;
}
}

View file

@ -0,0 +1,60 @@
package ru.spex.clearing.extdb.importer.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("extdb-importer")
public class ImportExtDBServiceSettings {
private HazelcastClientParams hazelcast;
private KafkaProducerSettings kafkaProducer;
private KafkaConsumerSettings kafkaConsumer;
private DatabaseSettings database;
private Cron cron;
public HazelcastClientParams getHazelcast() {
return hazelcast;
}
public void setHazelcast(HazelcastClientParams hazelcast) {
this.hazelcast = hazelcast;
}
public KafkaProducerSettings getKafkaProducer() {
return kafkaProducer;
}
public void setKafkaProducer(KafkaProducerSettings kafkaProducer) {
this.kafkaProducer = kafkaProducer;
}
public KafkaConsumerSettings getKafkaConsumer() {
return kafkaConsumer;
}
public void setKafkaConsumer(KafkaConsumerSettings kafkaConsumer) {
this.kafkaConsumer = kafkaConsumer;
}
public DatabaseSettings getDatabase() {
return database;
}
public void setDatabase(DatabaseSettings database) {
this.database = database;
}
public Cron getCron() {
return cron;
}
public void setCron(Cron cron) {
this.cron = cron;
}
}

View file

@ -0,0 +1,20 @@
package ru.spex.clearing.extdb.importer.error;
import ru.spcex.platform.utils.enumeration.IErrorEnumId;
public enum ExtDBImporterError implements IErrorEnumId {
crossRateNotValid(10002L),//todo код для crossRateNotValid??? Используется?
;
private final Long id;
ExtDBImporterError(Long id) {
this.id = id;
}
@Override
public Long getId() {
return id;
}
}

View file

@ -0,0 +1,19 @@
package ru.spex.clearing.extdb.importer.error;
public class ModuleInitializeException extends RuntimeException {
public ModuleInitializeException() {
}
public ModuleInitializeException(String message) {
super(message);
}
public ModuleInitializeException(String message, Throwable cause) {
super(message, cause);
}
public ModuleInitializeException(Throwable cause) {
super(cause);
}
}

View file

@ -0,0 +1,142 @@
package ru.spex.clearing.extdb.importer.services;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils;
import ru.clearing.classes.statics.data.misc.SCrossRate;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.STradesImportedRequest;
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.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IMessageResolver;
import java.math.BigDecimal;
import java.sql.Date;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Timestamp;
import java.time.Instant;
import java.time.LocalDate;
import java.util.Collection;
import java.util.HashMap;
import java.util.Map;
import java.util.function.Supplier;
import static ru.spex.clearing.extdb.importer.error.ExtDBImporterError.crossRateNotValid;
@Service
@EnableScheduling
public class ExtDBImporterService {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<SCrossRate> sCrossRateImdg;
private final JdbcTemplate jdbcTemplate;
private final Supplier<KafkaSender> kafka;
private final IMessageResolver messageResolver;
@Value("${extdb-importer.database.schema:dbo}")
private String schema;
public ExtDBImporterService(ImdgProvider imdgProvider, JdbcTemplate jdbcTemplate, Supplier<KafkaSender> kafka, IMessageResolver messageResolver) {
imdgProvider.waitAvailable();
this.sCrossRateImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SCrossRate, SCrossRate.class);
this.jdbcTemplate = jdbcTemplate;
this.kafka = kafka;
this.messageResolver = messageResolver;
}
@Scheduled(cron = "${extdb-importer.cron.load-from-db-cron}")
public synchronized void run() {
process(false);
}
public synchronized void process(boolean byCommand) {
log.debug("Start process import SCrossRate from DB. byCommand={}", byCommand);
String sql = String.format("SELECT * FROM %s.Trades", schema);
Collection<SCrossRate> crossRatesFromDB = jdbcTemplate.query(sql,
(resultSet, i) -> readSCrossRate(resultSet));
Instant currentInstant = Instant.now();
log.debug("load from DB SCrossRate \"{}\": {}", sql, crossRatesFromDB.size());
int created = 0;
int updated = 0;
for (SCrossRate crossRateDb : crossRatesFromDB) {
if (isValidCrossRate(crossRateDb)/* && fillNessessaryFields(crossRate)*/) {
SCrossRate sCrossRate = getSCrossRateFromImdg(crossRateDb, sCrossRateImdg);
if (sCrossRate != null && byCommand) {
// Обновление всех полей объекта из БД
crossRateDb.setId(sCrossRate.getId());
crossRateDb.setCreatedAt(sCrossRate.getCreatedAt());
crossRateDb.setUpdatedAt(currentInstant);
// В extDb нет таких полей, оставляем как есть:
crossRateDb.setGenerationTime(sCrossRate.getGenerationTime());
crossRateDb.setGenerationId(sCrossRate.getGenerationId());
sCrossRateImdg.update(crossRateDb);
updated++;
} else if (sCrossRate == null) {
// Создание объекта из БД
crossRateDb.setCreatedAt(currentInstant);
crossRateDb.setUpdatedAt(currentInstant);
sCrossRateImdg.insert(crossRateDb);
created++;
}
} else {
log.warn(messageResolver.resolve(new EnumMessage(crossRateNotValid, crossRateDb.getCurrCode())));
}
}
log.debug("Created {} new SCrossRate, {} updated.", created, updated);
if (created > 0 || updated > 0) {
if (byCommand) {
log.debug("Successfully import SCrossRate from DB by command.");
} else {
log.debug("Successfully import SCrossRate from DB by scheduled.");
}
// STradesImportedRequest importRequest = new STradesImportedRequest();
// kafka.get().sendRequestToQueue(Consts.S_CROSS_RATES_IMPORTED, importRequest);
} else {
log.debug("No SCrossRate were created or updated from DB. Kafka command will not be send");
}
}
public SCrossRate getSCrossRateFromImdg(SCrossRate crossRateDb, Imdg<SCrossRate> sCrossRateImdg) {
Map<String, Comparable<?>> query = new HashMap<>();
query.put("date", crossRateDb.getDate());
query.put("currCode", crossRateDb.getCurrCode());
query.put("currency", crossRateDb.getCurrency());
return sCrossRateImdg.getFirstObjectByFieldValues(query);
}
private boolean isValidCrossRate(SCrossRate crossRate) {
return crossRate.getDate() != null
&& (StringUtils.hasText(crossRate.getCurrCode())
|| StringUtils.hasText(crossRate.getCurrency()));
}
public SCrossRate readSCrossRate(ResultSet resultSet) throws SQLException {
SCrossRate object = new SCrossRate();
object.setDate(getLocalDateFromSqlDate(resultSet, "DATE"));
object.setCurrency(resultSet.getObject("CURRENCY", String.class));
object.setCurrCode(resultSet.getObject("CURR_CODE", String.class));
object.setFaceValue(resultSet.getObject("FACE_VALUE", BigDecimal.class));
object.setRate(resultSet.getObject("RATE", BigDecimal.class));
object.setUnitRate(resultSet.getObject("UNIT_RATE", BigDecimal.class));
return object;
}
private Instant getInstantFromTimestamp(ResultSet rs, String column) throws SQLException {
Timestamp date = rs.getTimestamp(column);
return date != null ? date.toInstant() : null;
}
private LocalDate getLocalDateFromSqlDate(ResultSet rs, String column) throws SQLException {
Date date = rs.getDate(column);
return date != null ? date.toLocalDate() : null;
}
}

View file

@ -0,0 +1,32 @@
package ru.spex.clearing.extdb.importer.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.beans.factory.annotation.Autowired;
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 ExtDBImporterService importer;
@Autowired
public LauncherCommandReceiver(Consumer<String, Object> kafkaQueue,
ExtDBImporterService importer) {
super(kafkaQueue);
this.importer = importer;
}
@Override
public void afterPropertiesSet() {
callback(LauncherCommandRequest.class)
.setConsumer(action -> importer.process(true))
.forDestination(Task.getCurrencyRates.topic(), callbacks::put); // CBRR
init();
}
}

View file

@ -0,0 +1,30 @@
spring.main.web-application-type=none
#extdb-importer.cron.load-from-db-cron=0 0/5 * * * ? - каждые 5 минут
extdb-importer.cron.load-from-db-cron=0 0/5 * * * ?
extdb-importer.database.login=CR_user
extdb-importer.database.password=Crossrate
extdb-importer.database.schema=dbo
extdb-importer.database.url=jdbc:sqlserver://10.200.200.144:1433;database=CBRInfo;
extdb-importer.database.driver=com.microsoft.sqlserver.jdbc.SQLServerDriver
extdb-importer.hazelcast.cluster-members=127.0.0.1:5701
extdb-importer.hazelcast.login=dev
extdb-importer.hazelcast.password=dev-pass
extdb-importer.kafka-consumer.bootstrap-servers=localhost:9092
extdb-importer.kafka-consumer.group-id=dev-group-extdb-importer
extdb-importer.kafka-consumer.enable-auto-commit=false
extdb-importer.kafka-consumer.session-timeout-ms=30000
extdb-importer.kafka-consumer.auto-offset-reset=latest
extdb-importer.kafka-consumer.linger-ms=1
extdb-importer.kafka-consumer.buffer-memory=33554432
extdb-importer.kafka-producer.bootstrap-servers=localhost:9092
extdb-importer.kafka-producer.acks=all
extdb-importer.kafka-producer.retries=0
extdb-importer.kafka-producer.batch-size=16384
extdb-importer.kafka-producer.linger-ms=1
extdb-importer.kafka-producer.buffer-memory=33554432

View file

@ -0,0 +1,56 @@
<?xml version="1.0" encoding="UTF-8"?>
<configuration>
<property name="LOG_PATH" value="./log" />
<property name="FILE_NAME" value="extdb-importer" />
<property name="CONSOLE_LOG_PATTERN" value="%date{HH:mm:ss.SSS} [%thread] %-5level %class{0}:%line - %message%n" />
<property name="FILE_LOG_PATTERN" value="%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %class{0}:%msg%n" />
<appender name="CONSOLE" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<!-- |%X{ru.nbch.scoring.web.logging.mdc_key}-->
<Pattern>${CONSOLE_LOG_PATTERN}</Pattern>
<charset>utf-8</charset>
</encoder>
</appender>
<!-- first FILE TEXT appender -->
<appender name="TEXT_FILE" class="ch.qos.logback.core.rolling.RollingFileAppender">
<file>${LOG_PATH}/${FILE_NAME}-text.log</file>
<encoder>
<!-- |%X{ru.nbch.scoring.web.logging.mdc_key}-->
<Pattern>${FILE_LOG_PATTERN}</Pattern>
<charset>utf8</charset>
</encoder>
<rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
<fileNamePattern>${LOG_PATH}/${FILE_NAME}-text.%d{yyyy-MM-dd}.%i.gz
</fileNamePattern>
<timeBasedFileNamingAndTriggeringPolicy class="ch.qos.logback.core.rolling.SizeAndTimeBasedFNATP">
<maxFileSize>100MB</maxFileSize>
</timeBasedFileNamingAndTriggeringPolicy>
<maxHistory>10</maxHistory>
</rollingPolicy>
</appender>
<!-- second FILE JSON appender -->
<appender name="JSON_FILE" class="ch.qos.logback.core.rolling.RollingFileAppender">
<file>${LOG_PATH}/${FILE_NAME}-json.log</file>
<encoder class="net.logstash.logback.encoder.LogstashEncoder" />
<rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
<fileNamePattern>${LOG_PATH}/${FILE_NAME}-json.%d{yyyy-MM-dd}.%i.gz
</fileNamePattern>
<timeBasedFileNamingAndTriggeringPolicy class="ch.qos.logback.core.rolling.SizeAndTimeBasedFNATP">
<maxFileSize>100MB</maxFileSize>
</timeBasedFileNamingAndTriggeringPolicy>
<maxHistory>10</maxHistory>
</rollingPolicy>
</appender>
<root level="info">
<!-- <appender-ref ref="CONSOLE"/> -->
<appender-ref ref="TEXT_FILE"/>
<appender-ref ref="JSON_FILE" />
</root>
<logger name="ru.spcex" level="debug" additivity="false">
<appender-ref ref="TEXT_FILE"/>
<appender-ref ref="JSON_FILE" />
<!-- <appender-ref ref="CONSOLE"/> -->
</logger>
</configuration>

View file

@ -0,0 +1,57 @@
package ru.spex.clearing.extdb.importer;
import com.hazelcast.core.Hazelcast;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.extension.ExtendWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.misc.SCrossRate;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.test.config.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spex.clearing.extdb.importer.config.ErrorResolverConfig;
import ru.spex.clearing.extdb.importer.config.DbTestConnectionConfig;
import ru.spex.clearing.extdb.importer.config.settings.ImportExtDBServiceSettings;
import ru.spex.clearing.extdb.importer.services.LauncherCommandReceiver;
import ru.spex.clearing.extdb.importer.services.ExtDBImporterService;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID;
import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId;
@ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = {
DbTestConnectionConfig.class,
ErrorResolverConfig.class,
ImportExtDBServiceSettings.class,
LauncherCommandReceiver.class,
ExtDBImporterService.class,
ImdgTestConfig.class,
KafkaTestConfig.class})
public abstract class AbstractServiceTest {
protected static final long id = currentID.getAndIncrement();
protected Imdg<SCrossRate> sCrossRateImdg;
@Autowired
@Qualifier("kafkaTestTemplate")
protected KafkaTemplate<String, Object> kafkaTemplate;
@Autowired
@Qualifier("hazelcastServiceTest")
protected ImdgProvider imdgProvider;
protected void init() {
waitAvailableImdgProviderAndAddAdminWithDefaultId();
this.sCrossRateImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SCrossRate, SCrossRate.class);
}
@AfterAll
static void shutdown() {
Hazelcast.shutdownAll();
}
}

View file

@ -0,0 +1,55 @@
package ru.spex.clearing.extdb.importer.config;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.jdbc.datasource.SingleConnectionDataSource;
import ru.spex.clearing.extdb.importer.error.ModuleInitializeException;
import javax.sql.DataSource;
import java.sql.Connection;
@SuppressWarnings("UnnecessaryLocalVariable")
@Configuration
public class DbTestConnectionConfig {
private final Logger log = LoggerFactory.getLogger(this.getClass());
@Bean(destroyMethod = "destroy")
public SingleConnectionDataSource dataSource() {
String login = "CR_user";
String password = "Crossrate";
String dbUrl = "jdbc:sqlserver://10.200.200.144:1433;database=CBRInfo";
SingleConnectionDataSource cpds = new SingleConnectionDataSource();
try {
cpds.setDriverClassName("com.microsoft.sqlserver.jdbc.SQLServerDriver");
} catch (Exception ue) {
throw new RuntimeException(ue);
}
cpds.setUrl(dbUrl);
cpds.setUsername(login);
cpds.setPassword(password);
String OPERATION_DATABASE_CONNECTION_CHECK = String.format("Database [%s] connection check", dbUrl);
try {
Connection conn = cpds.getConnection();
// conn.close();
log.info("{}: success", OPERATION_DATABASE_CONNECTION_CHECK);
return cpds;
} catch (Throwable e) {
String msg = String.format("%s: failed: %s -> %s",
OPERATION_DATABASE_CONNECTION_CHECK, e.getClass().getSimpleName(), e.getMessage());
log.error(msg);
throw new ModuleInitializeException(msg, e);
}
}
@Bean
public JdbcTemplate jdbcTemplate(DataSource dataSource) {
JdbcTemplate jdbcTemplate = new JdbcTemplate(dataSource);
return jdbcTemplate;
}
}

View file

@ -0,0 +1,108 @@
package ru.spex.clearing.extdb.importer.services;
import com.hazelcast.core.Hazelcast;
import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor;
import org.springframework.beans.factory.annotation.Autowired;
import ru.clearing.classes.statics.data.misc.SCrossRate;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.test.MatcherFactory;
import ru.spex.clearing.extdb.importer.AbstractServiceTest;
import ru.spcex.platform.enumeration.Task;
import javax.annotation.PostConstruct;
import java.math.BigDecimal;
import java.nio.file.Path;
import java.nio.file.Paths;
import java.time.LocalDate;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify;
import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.test.TestUtils.*;
import static ru.spcex.clearing.test.config.KafkaTestConfig.getCaptor;
class ExtDBImporterServiceTest extends AbstractServiceTest {
public static final MatcherFactory.Matcher<SCrossRate> S_CROSS_RATE_MATCHER = usingIgnoringFieldsComparator();
@Autowired
LauncherCommandReceiver launcherCommandReceiver;
@Autowired
ExtDBImporterService extDBImporterService;
@PostConstruct
public void init() {
super.init();
}
@BeforeAll
static void setProperty() {
Path path = Paths.get("src", "main", "resources");
String currentPath = path.toAbsolutePath().toString();
System.setProperty("spring.config.location", currentPath);
}
@BeforeEach
private void prepare(){
clearAllInImdg(sCrossRateImdg);
}
/**
* {@link ExtDBImporterService#process(boolean)}<br>
* Тест проверяет обновление сущности {@link SCrossRate}.<br>
* Входной запрос {@link LauncherCommandRequest}:<br>
*/
// @Test //для работы теста нужна тестовая база Microsoft SQL с данными
void process() throws InterruptedException {
SCrossRate sCrossRate = getSCrossRate();
Long id = sCrossRateImdg.insert(sCrossRate);
addRecordToKafka((MockConsumer) launcherCommandReceiver.getConsumer(), Task.getCurrencyRates.topic(), 0, 1, getJsonStringForNew(new LauncherCommandRequest(),0));
//waiting for kafka producer send message
// ArgumentCaptor<ProducerRecord> captor = getCaptor(kafkaTemplate);
// verify(kafkaTemplate, timeout(30_000L).times(1))
// .send(captor.capture());
Thread.sleep(1000L);
SCrossRate sCrossRates = extDBImporterService.getSCrossRateFromImdg(sCrossRate, sCrossRateImdg);
assertEquals(sCrossRates.getId(), id);
}
/**
* {@link ExtDBImporterService#process(boolean)}<br>
* Тест проверяет поиск сущности {@link SCrossRate}.<br>
* Входной запрос {@link LauncherCommandRequest}:<br>
*/
@Test
void getSCrossRateFromImdg(){
SCrossRate sCrossRate = getSCrossRate();
Long id = sCrossRateImdg.insert(sCrossRate);
SCrossRate sCrossRateRes = extDBImporterService.getSCrossRateFromImdg(sCrossRate, sCrossRateImdg);
S_CROSS_RATE_MATCHER.assertMatch(sCrossRateRes, sCrossRate);
sCrossRate.setCurrency("no-currency-now");
sCrossRateRes = extDBImporterService.getSCrossRateFromImdg(sCrossRate, sCrossRateImdg);
assertNull(sCrossRateRes);
}
private SCrossRate getSCrossRate(){
SCrossRate sCrossRate = new SCrossRate();
sCrossRate.setDate(LocalDate.of(2024,2,8));
sCrossRate.setRate(new BigDecimal("10.1"));
sCrossRate.setUnitRate(new BigDecimal("12.4"));
sCrossRate.setCurrCode("RUB");
sCrossRate.setCurrency("RUB");
return sCrossRate;
}
}

View file

@ -37,6 +37,7 @@
<module>test-clearing</module>
<module>cleaning-builders</module>
<module>trade-importer</module>
<module>extdb-importer</module>
<module>lim-exporter</module>
<module>swt-exporter</module>
<module>swt-importer</module>

View file

@ -7,6 +7,7 @@ public enum Task implements IEnumKey {
accountBlock("ABLK"),//Блокировка счета
additionOrDeleteOfBalance("ADBL"),//Дозачисление/списание остатков
getOfTrades("GTRD"),//Получение сделок из Торговой системы
getCurrencyRates("CBRR"),//Загрузка кросс-курсов
getVerification("GVER"),// Запуск сверки
@Deprecated /* todo GBLD удаляется по CLS-267, CLS-275 */ getBalance("GBLD"),// Поступление средств
startOfClearing("SCLR"),// Запуск клиринговой сессии

View file

@ -42,6 +42,7 @@
<folder_root_swt-exporter>${folder_root_clearing}/clearing-parent/swt-exporter</folder_root_swt-exporter>
<folder_root_swt-importer>${folder_root_clearing}/clearing-parent/swt-importer</folder_root_swt-importer>
<folder_root_trade-importer>${folder_root_clearing}/clearing-parent/trade-importer</folder_root_trade-importer>
<folder_root_extdb-importer>${folder_root_clearing}/clearing-parent/extdb-importer</folder_root_extdb-importer>
<folder_root_account-service>${folder_root_clearing}/clearing-parent/account-service</folder_root_account-service>
<folder_root_balance-service>${folder_root_clearing}/clearing-parent/balance-service</folder_root_balance-service>
<folder_root_company-service>${folder_root_clearing}/clearing-parent/company-service</folder_root_company-service>

View file

@ -293,6 +293,25 @@
</fileSets>
</configuration>
</execution>
<execution>
<id>copy-extdb-importer-bin</id>
<phase>prepare-package</phase>
<goals>
<goal>copy</goal>
</goals>
<configuration>
<fileSets>
<fileSet>
<sourceFile>${folder_root_extdb-importer}/target/extdb-importer.jar</sourceFile>
<destinationFile>${folder.clearing.distr.bin}/extdb-importer.jar</destinationFile>
</fileSet>
<fileSet>
<sourceFile>${folder_root_extdb-importer}/src/main/resources/application.properties</sourceFile>
<destinationFile>${folder.clearing.distr.settings}/extdb-importer/application.properties</destinationFile>
</fileSet>
</fileSets>
</configuration>
</execution>
<execution>
<id>copy-account-service-bin</id>
<phase>prepare-package</phase>

View file

@ -0,0 +1,9 @@
#!/bin/bash
CLEARING_HOME=/opt/mfd/clearing/
cd $CLEARING_HOME/bin
CMD="java -Xrunjdwp:transport=dt_socket,server=y,suspend=n,address=7100 -jar trade-importer.jar --spring.config.location=$CLEARING_HOME/settings/trade-importer/"
$CMD >/dev/null 2>&1 &

View file

@ -12,6 +12,7 @@ cd /opt/mfd/clearing/bin
/opt/mfd/clearing/bin/swt-exporter.sh
/opt/mfd/clearing/bin/dbf-importer.sh
/opt/mfd/clearing/bin/trade-importer.sh
/opt/mfd/clearing/bin/extdb-importer.sh
/opt/mfd/clearing/bin/securities-service.sh
/opt/mfd/clearing/bin/utility-service.sh
/opt/mfd/clearing/bin/scheduler-service.sh

View file

@ -3,7 +3,7 @@
CLEARING_HOME=/opt/mfd/clearing/
cd $CLEARING_HOME/bin
CMD="java -Xrunjdwp:transport=dt_socket,server=y,suspend=n,address=7100 -jar trade-importer.jar --spring.config.location=$CLEARING_HOME/settings/trade-importer/"
CMD="java -Xrunjdwp:transport=dt_socket,server=y,suspend=n,address=7140 -jar extdb-importer.jar --spring.config.location=$CLEARING_HOME/settings/extdb-importer/"
$CMD >/dev/null 2>&1 &