parent
7e4f7788e6
commit
b5592ad393
10 changed files with 131 additions and 35 deletions
|
|
@ -23,6 +23,11 @@ public class STradesMapStore extends TemplateMapStore<STrades> {
|
||||||
return IMDGDistributedNames.Map_STrades;
|
return IMDGDistributedNames.Map_STrades;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public String[] getIndexingField() {
|
||||||
|
return new String[]{"tradeNum"};// список индексируемых полей
|
||||||
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public String getTableName() {
|
public String getTableName() {
|
||||||
return "S_TRADES";
|
return "S_TRADES";
|
||||||
|
|
|
||||||
|
|
@ -13,6 +13,8 @@ import ru.spcex.clearing.trade.importer.error.ModuleInitializeException;
|
||||||
import javax.sql.DataSource;
|
import javax.sql.DataSource;
|
||||||
import java.sql.Connection;
|
import java.sql.Connection;
|
||||||
|
|
||||||
|
import static org.springframework.jdbc.datasource.DataSourceUtils.doCloseConnection;
|
||||||
|
|
||||||
@SuppressWarnings("UnnecessaryLocalVariable")
|
@SuppressWarnings("UnnecessaryLocalVariable")
|
||||||
@Configuration
|
@Configuration
|
||||||
public class DbConnectionConfig {
|
public class DbConnectionConfig {
|
||||||
|
|
@ -29,23 +31,24 @@ public class DbConnectionConfig {
|
||||||
String login = settings.getLogin();
|
String login = settings.getLogin();
|
||||||
String password = settings.getPassword();
|
String password = settings.getPassword();
|
||||||
String dbUrl = settings.getUrl();
|
String dbUrl = settings.getUrl();
|
||||||
|
String driver = settings.getDriver();
|
||||||
|
|
||||||
SingleConnectionDataSource cpds = new SingleConnectionDataSource();
|
SingleConnectionDataSource ds = new SingleConnectionDataSource();
|
||||||
try {
|
try {
|
||||||
cpds.setDriverClassName("com.microsoft.sqlserver.jdbc.SQLServerDriver");
|
ds.setDriverClassName(driver);
|
||||||
} catch (Exception ue) {
|
} catch (Exception ue) {
|
||||||
throw new RuntimeException(ue);
|
throw new RuntimeException(ue);
|
||||||
}
|
}
|
||||||
cpds.setUrl(dbUrl);
|
ds.setUrl(dbUrl);
|
||||||
cpds.setUsername(login);
|
ds.setUsername(login);
|
||||||
cpds.setPassword(password);
|
ds.setPassword(password);
|
||||||
|
|
||||||
String OPERATION_DATABASE_CONNECTION_CHECK = String.format("Database [%s] connection check", dbUrl);
|
String OPERATION_DATABASE_CONNECTION_CHECK = String.format("Database [%s] connection check", dbUrl);
|
||||||
try {
|
try {
|
||||||
Connection conn = cpds.getConnection();
|
Connection conn = ds.getConnection();
|
||||||
// conn.close();
|
doCloseConnection(conn, ds);
|
||||||
log.info("{}: success", OPERATION_DATABASE_CONNECTION_CHECK);
|
log.info("{}: success", OPERATION_DATABASE_CONNECTION_CHECK);
|
||||||
return cpds;
|
return ds;
|
||||||
} catch (Throwable e) {
|
} catch (Throwable e) {
|
||||||
String msg = String.format("%s: failed: %s -> %s",
|
String msg = String.format("%s: failed: %s -> %s",
|
||||||
OPERATION_DATABASE_CONNECTION_CHECK, e.getClass().getSimpleName(), e.getMessage());
|
OPERATION_DATABASE_CONNECTION_CHECK, e.getClass().getSimpleName(), e.getMessage());
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,14 @@
|
||||||
|
package ru.spcex.clearing.trade.importer.config;
|
||||||
|
|
||||||
|
import org.springframework.context.annotation.Bean;
|
||||||
|
import org.springframework.context.annotation.Configuration;
|
||||||
|
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
||||||
|
import ru.spcex.platform.utils.enumeration.SimpleMessageResolver;
|
||||||
|
|
||||||
|
@Configuration
|
||||||
|
public class ErrorResolverConfig {
|
||||||
|
@Bean
|
||||||
|
public IMessageResolver messageResolver() {
|
||||||
|
return new SimpleMessageResolver();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -4,6 +4,7 @@ public class DatabaseSettings {
|
||||||
private String login;
|
private String login;
|
||||||
private String password;
|
private String password;
|
||||||
private String url;
|
private String url;
|
||||||
|
private String driver;
|
||||||
|
|
||||||
public String getLogin() {
|
public String getLogin() {
|
||||||
return login;
|
return login;
|
||||||
|
|
@ -28,4 +29,12 @@ public class DatabaseSettings {
|
||||||
public void setUrl(String url) {
|
public void setUrl(String url) {
|
||||||
this.url = url;
|
this.url = url;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public String getDriver() {
|
||||||
|
return driver;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setDriver(String driver) {
|
||||||
|
this.driver = driver;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,20 @@
|
||||||
|
package ru.spcex.clearing.trade.importer.error;
|
||||||
|
|
||||||
|
import ru.spcex.platform.utils.enumeration.IErrorEnumId;
|
||||||
|
|
||||||
|
public enum ValidationError implements IErrorEnumId {
|
||||||
|
sTradesNotValid(10001L),
|
||||||
|
;
|
||||||
|
|
||||||
|
private final Long id;
|
||||||
|
|
||||||
|
ValidationError(Long id) {
|
||||||
|
this.id = id;
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public Long getId() {
|
||||||
|
return id;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -25,7 +25,7 @@ public class LauncherCommandReceiver extends QueueConsumer implements Initializi
|
||||||
@Override
|
@Override
|
||||||
public void afterPropertiesSet() {
|
public void afterPropertiesSet() {
|
||||||
callback(LauncherCommandRequest.class)
|
callback(LauncherCommandRequest.class)
|
||||||
.setConsumer(action -> importer.process())
|
.setConsumer(action -> importer.process(true))
|
||||||
.forDestination(Task.getOfTrades.topic(), callbacks::put); // GTRD
|
.forDestination(Task.getOfTrades.topic(), callbacks::put); // GTRD
|
||||||
init();
|
init();
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -16,11 +16,14 @@ import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.STradesImportedRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.STradesImportedRequest;
|
||||||
import ru.spcex.platform.imdg.api.Imdg;
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
|
import ru.spcex.platform.utils.enumeration.EnumMessage;
|
||||||
|
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
||||||
|
|
||||||
import java.util.Collection;
|
import java.util.Collection;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
|
|
||||||
import static ru.spcex.clearing.platform.messaging.domain.Consts.S_TRADES_IMPORTED;
|
import static ru.spcex.clearing.platform.messaging.domain.Consts.S_TRADES_IMPORTED;
|
||||||
|
import static ru.spcex.clearing.trade.importer.error.ValidationError.sTradesNotValid;
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
@EnableScheduling
|
@EnableScheduling
|
||||||
|
|
@ -29,38 +32,46 @@ public class TradeImporterService {
|
||||||
private final Imdg<STrades> sTradesImdg;
|
private final Imdg<STrades> sTradesImdg;
|
||||||
private final JdbcTemplate jdbcTemplate;
|
private final JdbcTemplate jdbcTemplate;
|
||||||
private final Producer<String, Object> producer;
|
private final Producer<String, Object> producer;
|
||||||
|
private final IMessageResolver messageResolver;
|
||||||
private static final RowMapper<STrades> ROW_MAPPER = BeanPropertyRowMapper.newInstance(STrades.class);
|
private static final RowMapper<STrades> ROW_MAPPER = BeanPropertyRowMapper.newInstance(STrades.class);
|
||||||
|
|
||||||
|
|
||||||
public TradeImporterService(ImdgProvider imdgProvider, JdbcTemplate jdbcTemplate, Producer<String, Object> producer) {
|
public TradeImporterService(ImdgProvider imdgProvider, JdbcTemplate jdbcTemplate, Producer<String, Object> producer, IMessageResolver messageResolver) {
|
||||||
this.sTradesImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_STrades, STrades.class);
|
this.sTradesImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_STrades, STrades.class);
|
||||||
this.jdbcTemplate = jdbcTemplate;
|
this.jdbcTemplate = jdbcTemplate;
|
||||||
this.producer = producer;
|
this.producer = producer;
|
||||||
|
this.messageResolver = messageResolver;
|
||||||
}
|
}
|
||||||
|
|
||||||
@Scheduled(cron = "${trade-importer.cron.load-from-db-cron}")
|
@Scheduled(cron = "${trade-importer.cron.load-from-db-cron}")
|
||||||
public void run() {
|
public void run() {
|
||||||
process();
|
process(false);
|
||||||
}
|
}
|
||||||
|
|
||||||
public void process() {
|
public void process(boolean byCommand) {
|
||||||
Collection<STrades> tradesFromDB = jdbcTemplate.query("SELECT * FROM Trades", ROW_MAPPER);
|
Collection<STrades> tradesFromDB = jdbcTemplate.query("SELECT * FROM Trades", ROW_MAPPER);
|
||||||
|
|
||||||
for (STrades tradesDb : tradesFromDB) {
|
for (STrades tradesDb : tradesFromDB) {
|
||||||
if (isValidTrades(tradesDb)) {
|
if (isValidTrades(tradesDb)) {
|
||||||
STrades sTrades = sTradesImdg.getSingleObjectByFieldValues(Map.of("tradeDate", tradesDb.getTradeDate(),
|
STrades sTrades = getSTradesFromImdg(tradesDb, sTradesImdg);
|
||||||
"tradeNum", tradesDb.getTradeNum(),
|
|
||||||
"operation", tradesDb.getOperation(),
|
|
||||||
"classCode", tradesDb.getClassCode()));
|
|
||||||
if (sTrades != null) tradesDb.setId(sTrades.getId());
|
if (sTrades != null) tradesDb.setId(sTrades.getId());
|
||||||
sTradesImdg.insert(tradesDb);
|
sTradesImdg.insert(tradesDb);
|
||||||
} else {
|
} else {
|
||||||
log.warn(String.format("(10001) \"Сделка с номером в ТС = %s некорректна\"", tradesDb.getTradeNum()));
|
log.warn(messageResolver.resolve(new EnumMessage(sTradesNotValid, tradesDb.getTradeNum())));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
STradesImportedRequest sTradesImportedRequest = new STradesImportedRequest();
|
if (byCommand) {
|
||||||
producer.send(new ProducerRecord<>(S_TRADES_IMPORTED, sTradesImportedRequest));
|
STradesImportedRequest sTradesImportedRequest = new STradesImportedRequest();
|
||||||
|
producer.send(new ProducerRecord<>(S_TRADES_IMPORTED, sTradesImportedRequest));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
public STrades getSTradesFromImdg(STrades tradesDb, Imdg<STrades> sTradesImdg) {
|
||||||
|
return sTradesImdg.getSingleObjectByFieldValues(Map.of("tradeDate", tradesDb.getTradeDate(),
|
||||||
|
"tradeNum", tradesDb.getTradeNum(),
|
||||||
|
"operation", tradesDb.getOperation(),
|
||||||
|
"classCode", tradesDb.getClassCode()));
|
||||||
}
|
}
|
||||||
|
|
||||||
private boolean isValidTrades(STrades trades) {
|
private boolean isValidTrades(STrades trades) {
|
||||||
|
|
|
||||||
|
|
@ -6,6 +6,7 @@ trade-importer.cron.load-from-db-cron=0 0/5 * * * ?
|
||||||
trade-importer.database.login=sa
|
trade-importer.database.login=sa
|
||||||
trade-importer.database.password=Aa123456
|
trade-importer.database.password=Aa123456
|
||||||
trade-importer.database.url=jdbc:sqlserver://localhost:1433;database=SPVB_TS;schema=dbo
|
trade-importer.database.url=jdbc:sqlserver://localhost:1433;database=SPVB_TS;schema=dbo
|
||||||
|
trade-importer.database.driver=com.microsoft.sqlserver.jdbc.SQLServerDriver
|
||||||
|
|
||||||
trade-importer.hazelcast.cluster-members=127.0.0.1:5701
|
trade-importer.hazelcast.cluster-members=127.0.0.1:5701
|
||||||
trade-importer.hazelcast.login=dev
|
trade-importer.hazelcast.login=dev
|
||||||
|
|
|
||||||
|
|
@ -17,6 +17,7 @@ import ru.spcex.clearing.test.MatcherFactory;
|
||||||
import ru.spcex.clearing.test.TestUtils;
|
import ru.spcex.clearing.test.TestUtils;
|
||||||
import ru.spcex.clearing.test.config.ImdgTestConfig;
|
import ru.spcex.clearing.test.config.ImdgTestConfig;
|
||||||
import ru.spcex.clearing.test.config.KafkaTestConfig;
|
import ru.spcex.clearing.test.config.KafkaTestConfig;
|
||||||
|
import ru.spcex.clearing.trade.importer.config.ErrorResolverConfig;
|
||||||
import ru.spcex.clearing.trade.importer.config.config.DbTestConnectionConfig;
|
import ru.spcex.clearing.trade.importer.config.config.DbTestConnectionConfig;
|
||||||
import ru.spcex.clearing.trade.importer.config.settings.ImportTradeServiceSettings;
|
import ru.spcex.clearing.trade.importer.config.settings.ImportTradeServiceSettings;
|
||||||
import ru.spcex.clearing.trade.importer.services.LauncherCommandReceiver;
|
import ru.spcex.clearing.trade.importer.services.LauncherCommandReceiver;
|
||||||
|
|
@ -33,6 +34,7 @@ import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProv
|
||||||
@ExtendWith(SpringExtension.class)
|
@ExtendWith(SpringExtension.class)
|
||||||
@ContextConfiguration(classes = {
|
@ContextConfiguration(classes = {
|
||||||
DbTestConnectionConfig.class,
|
DbTestConnectionConfig.class,
|
||||||
|
ErrorResolverConfig.class,
|
||||||
ImportTradeServiceSettings.class,
|
ImportTradeServiceSettings.class,
|
||||||
LauncherCommandReceiver.class,
|
LauncherCommandReceiver.class,
|
||||||
TradeImporterService.class,
|
TradeImporterService.class,
|
||||||
|
|
|
||||||
|
|
@ -2,9 +2,12 @@ package ru.spcex.clearing.trade.importer.services;
|
||||||
|
|
||||||
import org.apache.kafka.clients.consumer.MockConsumer;
|
import org.apache.kafka.clients.consumer.MockConsumer;
|
||||||
import org.junit.jupiter.api.BeforeAll;
|
import org.junit.jupiter.api.BeforeAll;
|
||||||
|
import org.junit.jupiter.api.BeforeEach;
|
||||||
|
import org.junit.jupiter.api.Test;
|
||||||
import org.springframework.beans.factory.annotation.Autowired;
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
import ru.clearing.classes.statics.data.misc.STrades;
|
import ru.clearing.classes.statics.data.misc.STrades;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
|
||||||
|
import ru.spcex.clearing.test.MatcherFactory;
|
||||||
import ru.spcex.clearing.trade.importer.AbstractServiceTest;
|
import ru.spcex.clearing.trade.importer.AbstractServiceTest;
|
||||||
import ru.spcex.platform.enumeration.Task;
|
import ru.spcex.platform.enumeration.Task;
|
||||||
|
|
||||||
|
|
@ -12,18 +15,23 @@ import javax.annotation.PostConstruct;
|
||||||
import java.nio.file.Path;
|
import java.nio.file.Path;
|
||||||
import java.nio.file.Paths;
|
import java.nio.file.Paths;
|
||||||
import java.time.LocalDate;
|
import java.time.LocalDate;
|
||||||
import java.util.Map;
|
|
||||||
|
|
||||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
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.timeout;
|
||||||
import static org.mockito.Mockito.verify;
|
import static org.mockito.Mockito.verify;
|
||||||
import static ru.spcex.clearing.test.TestUtils.addRecordToKafka;
|
import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator;
|
||||||
import static ru.spcex.clearing.test.TestUtils.getJsonStringForNew;
|
import static ru.spcex.clearing.test.TestUtils.*;
|
||||||
|
|
||||||
class TradeImporterServiceTest extends AbstractServiceTest {
|
class TradeImporterServiceTest extends AbstractServiceTest {
|
||||||
|
public static final MatcherFactory.Matcher<STrades> S_TRADES_MATCHER = usingIgnoringFieldsComparator();
|
||||||
|
|
||||||
@Autowired
|
@Autowired
|
||||||
LauncherCommandReceiver launcherCommandReceiver;
|
LauncherCommandReceiver launcherCommandReceiver;
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
TradeImporterService tradeImporterService;
|
||||||
|
|
||||||
@PostConstruct
|
@PostConstruct
|
||||||
public void init() {
|
public void init() {
|
||||||
super.init();
|
super.init();
|
||||||
|
|
@ -37,20 +45,20 @@ class TradeImporterServiceTest extends AbstractServiceTest {
|
||||||
// Hazelcast.shutdownAll();
|
// Hazelcast.shutdownAll();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@BeforeEach
|
||||||
|
private void prepare(){
|
||||||
|
clearAllInImdg(sTradesImdg);
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* {@link TradeImporterService#process()}<br>
|
* {@link TradeImporterService#process(boolean)}<br>
|
||||||
* Тест проверяет обновление сущности {@link STrades}.<br>
|
* Тест проверяет обновление сущности {@link STrades}.<br>
|
||||||
* Входной запрос {@link LauncherCommandRequest}:<br>
|
* Входной запрос {@link LauncherCommandRequest}:<br>
|
||||||
*/
|
*/
|
||||||
// @Test для работы теста нужна тестовая база Microsoft SQL с данными
|
@Test //для работы теста нужна тестовая база Microsoft SQL с данными
|
||||||
void process() {
|
void process() {
|
||||||
STrades trades = new STrades();
|
STrades sTrade = getSTrade();
|
||||||
trades.setId(22L);
|
Long id = sTradesImdg.insert(sTrade);
|
||||||
trades.setTradeDate(LocalDate.of(2023,4,19));
|
|
||||||
trades.setTradeNum(661486L);
|
|
||||||
trades.setOperation("operation20");
|
|
||||||
trades.setClassCode("UESC");
|
|
||||||
Long id = sTradesImdg.insert(trades);
|
|
||||||
|
|
||||||
addRecordToKafka((MockConsumer) launcherCommandReceiver.getConsumer(), Task.getOfTrades.topic(), 0, 1, getJsonStringForNew(new LauncherCommandRequest(),0));
|
addRecordToKafka((MockConsumer) launcherCommandReceiver.getConsumer(), Task.getOfTrades.topic(), 0, 1, getJsonStringForNew(new LauncherCommandRequest(),0));
|
||||||
|
|
||||||
|
|
@ -58,11 +66,34 @@ class TradeImporterServiceTest extends AbstractServiceTest {
|
||||||
verify(mockProducer, timeout(30_000L).times(1))
|
verify(mockProducer, timeout(30_000L).times(1))
|
||||||
.send(producerRecord.capture());
|
.send(producerRecord.capture());
|
||||||
|
|
||||||
STrades sTrades = sTradesImdg.getSingleObjectByFieldValues(Map.of("tradeDate", trades.getTradeDate(),
|
STrades sTrades = tradeImporterService.getSTradesFromImdg(sTrade, sTradesImdg);
|
||||||
"tradeNum", trades.getTradeNum(),
|
|
||||||
"operation", trades.getOperation(),
|
|
||||||
"classCode", trades.getClassCode()));
|
|
||||||
|
|
||||||
assertEquals(sTrades.getId(), id);
|
assertEquals(sTrades.getId(), id);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* {@link TradeImporterService#process(boolean)}<br>
|
||||||
|
* Тест проверяет обновление сущности {@link STrades}.<br>
|
||||||
|
* Входной запрос {@link LauncherCommandRequest}:<br>
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void getSTradesFromImdg(){
|
||||||
|
STrades sTrade = getSTrade();
|
||||||
|
Long id = sTradesImdg.insert(sTrade);
|
||||||
|
STrades sTradeRes = tradeImporterService.getSTradesFromImdg(sTrade, sTradesImdg);
|
||||||
|
S_TRADES_MATCHER.assertMatch(sTradeRes, sTrade);
|
||||||
|
|
||||||
|
sTrade.setTradeNum(23L);
|
||||||
|
sTradeRes = tradeImporterService.getSTradesFromImdg(sTrade, sTradesImdg);
|
||||||
|
assertNull(sTradeRes);
|
||||||
|
}
|
||||||
|
|
||||||
|
private STrades getSTrade(){
|
||||||
|
STrades trades = new STrades();
|
||||||
|
trades.setTradeDate(LocalDate.of(2023,4,19));
|
||||||
|
trades.setTradeNum(661486L);
|
||||||
|
trades.setOperation("operation20");
|
||||||
|
trades.setClassCode("UESC");
|
||||||
|
return trades;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
Loading…
Add table
Reference in a new issue