diff --git a/clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/object/STradesMapStore.java b/clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/object/STradesMapStore.java index 22f0b37b1..1c7c4a923 100644 --- a/clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/object/STradesMapStore.java +++ b/clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/object/STradesMapStore.java @@ -23,6 +23,11 @@ public class STradesMapStore extends TemplateMapStore { return IMDGDistributedNames.Map_STrades; } + @Override + public String[] getIndexingField() { + return new String[]{"tradeNum"};// список индексируемых полей + } + @Override public String getTableName() { return "S_TRADES"; diff --git a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/config/DbConnectionConfig.java b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/config/DbConnectionConfig.java index 3afddae6b..103dd0497 100644 --- a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/config/DbConnectionConfig.java +++ b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/config/DbConnectionConfig.java @@ -13,6 +13,8 @@ import ru.spcex.clearing.trade.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 { @@ -29,23 +31,24 @@ public class DbConnectionConfig { String login = settings.getLogin(); String password = settings.getPassword(); String dbUrl = settings.getUrl(); + String driver = settings.getDriver(); - SingleConnectionDataSource cpds = new SingleConnectionDataSource(); + SingleConnectionDataSource ds = new SingleConnectionDataSource(); try { - cpds.setDriverClassName("com.microsoft.sqlserver.jdbc.SQLServerDriver"); + ds.setDriverClassName(driver); } catch (Exception ue) { throw new RuntimeException(ue); } - cpds.setUrl(dbUrl); - cpds.setUsername(login); - cpds.setPassword(password); + ds.setUrl(dbUrl); + ds.setUsername(login); + ds.setPassword(password); String OPERATION_DATABASE_CONNECTION_CHECK = String.format("Database [%s] connection check", dbUrl); try { - Connection conn = cpds.getConnection(); -// conn.close(); + Connection conn = ds.getConnection(); + doCloseConnection(conn, ds); log.info("{}: success", OPERATION_DATABASE_CONNECTION_CHECK); - return cpds; + return ds; } catch (Throwable e) { String msg = String.format("%s: failed: %s -> %s", OPERATION_DATABASE_CONNECTION_CHECK, e.getClass().getSimpleName(), e.getMessage()); diff --git a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/config/ErrorResolverConfig.java b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/config/ErrorResolverConfig.java new file mode 100644 index 000000000..8c838d4ed --- /dev/null +++ b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/config/ErrorResolverConfig.java @@ -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(); + } +} diff --git a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/config/settings/DatabaseSettings.java b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/config/settings/DatabaseSettings.java index c019613e1..1c84d9ec7 100644 --- a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/config/settings/DatabaseSettings.java +++ b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/config/settings/DatabaseSettings.java @@ -4,6 +4,7 @@ public class DatabaseSettings { private String login; private String password; private String url; + private String driver; public String getLogin() { return login; @@ -28,4 +29,12 @@ public class DatabaseSettings { public void setUrl(String url) { this.url = url; } + + public String getDriver() { + return driver; + } + + public void setDriver(String driver) { + this.driver = driver; + } } diff --git a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/error/ValidationError.java b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/error/ValidationError.java new file mode 100644 index 000000000..dbbbbb8c3 --- /dev/null +++ b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/error/ValidationError.java @@ -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; + } +} diff --git a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/LauncherCommandReceiver.java b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/LauncherCommandReceiver.java index 074343102..0d92824f2 100644 --- a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/LauncherCommandReceiver.java +++ b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/LauncherCommandReceiver.java @@ -25,7 +25,7 @@ public class LauncherCommandReceiver extends QueueConsumer implements Initializi @Override public void afterPropertiesSet() { callback(LauncherCommandRequest.class) - .setConsumer(action -> importer.process()) + .setConsumer(action -> importer.process(true)) .forDestination(Task.getOfTrades.topic(), callbacks::put); // GTRD init(); } diff --git a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/TradeImporterService.java b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/TradeImporterService.java index 22a023f7f..08f00d9cf 100644 --- a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/TradeImporterService.java +++ b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/TradeImporterService.java @@ -16,11 +16,14 @@ import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.domain.cud.utilities.STradesImportedRequest; 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.util.Collection; import java.util.Map; import static ru.spcex.clearing.platform.messaging.domain.Consts.S_TRADES_IMPORTED; +import static ru.spcex.clearing.trade.importer.error.ValidationError.sTradesNotValid; @Service @EnableScheduling @@ -29,38 +32,46 @@ public class TradeImporterService { private final Imdg sTradesImdg; private final JdbcTemplate jdbcTemplate; private final Producer producer; + private final IMessageResolver messageResolver; private static final RowMapper ROW_MAPPER = BeanPropertyRowMapper.newInstance(STrades.class); - public TradeImporterService(ImdgProvider imdgProvider, JdbcTemplate jdbcTemplate, Producer producer) { + public TradeImporterService(ImdgProvider imdgProvider, JdbcTemplate jdbcTemplate, Producer producer, IMessageResolver messageResolver) { this.sTradesImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_STrades, STrades.class); this.jdbcTemplate = jdbcTemplate; this.producer = producer; + this.messageResolver = messageResolver; } @Scheduled(cron = "${trade-importer.cron.load-from-db-cron}") public void run() { - process(); + process(false); } - public void process() { + public void process(boolean byCommand) { Collection tradesFromDB = jdbcTemplate.query("SELECT * FROM Trades", ROW_MAPPER); for (STrades tradesDb : tradesFromDB) { if (isValidTrades(tradesDb)) { - STrades sTrades = sTradesImdg.getSingleObjectByFieldValues(Map.of("tradeDate", tradesDb.getTradeDate(), - "tradeNum", tradesDb.getTradeNum(), - "operation", tradesDb.getOperation(), - "classCode", tradesDb.getClassCode())); + STrades sTrades = getSTradesFromImdg(tradesDb, sTradesImdg); if (sTrades != null) tradesDb.setId(sTrades.getId()); sTradesImdg.insert(tradesDb); } else { - log.warn(String.format("(10001) \"Сделка с номером в ТС = %s некорректна\"", tradesDb.getTradeNum())); + log.warn(messageResolver.resolve(new EnumMessage(sTradesNotValid, tradesDb.getTradeNum()))); } } - STradesImportedRequest sTradesImportedRequest = new STradesImportedRequest(); - producer.send(new ProducerRecord<>(S_TRADES_IMPORTED, sTradesImportedRequest)); + if (byCommand) { + STradesImportedRequest sTradesImportedRequest = new STradesImportedRequest(); + producer.send(new ProducerRecord<>(S_TRADES_IMPORTED, sTradesImportedRequest)); + } + } + + public STrades getSTradesFromImdg(STrades tradesDb, Imdg sTradesImdg) { + return sTradesImdg.getSingleObjectByFieldValues(Map.of("tradeDate", tradesDb.getTradeDate(), + "tradeNum", tradesDb.getTradeNum(), + "operation", tradesDb.getOperation(), + "classCode", tradesDb.getClassCode())); } private boolean isValidTrades(STrades trades) { diff --git a/clearing-parent/trade-importer/src/main/resources/application.properties b/clearing-parent/trade-importer/src/main/resources/application.properties index fb9ca6cef..dc8f7ba47 100644 --- a/clearing-parent/trade-importer/src/main/resources/application.properties +++ b/clearing-parent/trade-importer/src/main/resources/application.properties @@ -6,6 +6,7 @@ trade-importer.cron.load-from-db-cron=0 0/5 * * * ? trade-importer.database.login=sa trade-importer.database.password=Aa123456 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.login=dev diff --git a/clearing-parent/trade-importer/src/test/java/ru/spcex/clearing/trade/importer/AbstractServiceTest.java b/clearing-parent/trade-importer/src/test/java/ru/spcex/clearing/trade/importer/AbstractServiceTest.java index 3e1e33c3d..84949d3e3 100644 --- a/clearing-parent/trade-importer/src/test/java/ru/spcex/clearing/trade/importer/AbstractServiceTest.java +++ b/clearing-parent/trade-importer/src/test/java/ru/spcex/clearing/trade/importer/AbstractServiceTest.java @@ -17,6 +17,7 @@ import ru.spcex.clearing.test.MatcherFactory; import ru.spcex.clearing.test.TestUtils; import ru.spcex.clearing.test.config.ImdgTestConfig; 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.settings.ImportTradeServiceSettings; import ru.spcex.clearing.trade.importer.services.LauncherCommandReceiver; @@ -33,6 +34,7 @@ import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProv @ExtendWith(SpringExtension.class) @ContextConfiguration(classes = { DbTestConnectionConfig.class, + ErrorResolverConfig.class, ImportTradeServiceSettings.class, LauncherCommandReceiver.class, TradeImporterService.class, diff --git a/clearing-parent/trade-importer/src/test/java/ru/spcex/clearing/trade/importer/services/TradeImporterServiceTest.java b/clearing-parent/trade-importer/src/test/java/ru/spcex/clearing/trade/importer/services/TradeImporterServiceTest.java index 615339e94..49068bcb9 100644 --- a/clearing-parent/trade-importer/src/test/java/ru/spcex/clearing/trade/importer/services/TradeImporterServiceTest.java +++ b/clearing-parent/trade-importer/src/test/java/ru/spcex/clearing/trade/importer/services/TradeImporterServiceTest.java @@ -2,9 +2,12 @@ package ru.spcex.clearing.trade.importer.services; import org.apache.kafka.clients.consumer.MockConsumer; 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 ru.clearing.classes.statics.data.misc.STrades; 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.platform.enumeration.Task; @@ -12,18 +15,23 @@ import javax.annotation.PostConstruct; import java.nio.file.Path; import java.nio.file.Paths; import java.time.LocalDate; -import java.util.Map; 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.TestUtils.addRecordToKafka; -import static ru.spcex.clearing.test.TestUtils.getJsonStringForNew; +import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator; +import static ru.spcex.clearing.test.TestUtils.*; class TradeImporterServiceTest extends AbstractServiceTest { + public static final MatcherFactory.Matcher S_TRADES_MATCHER = usingIgnoringFieldsComparator(); + @Autowired LauncherCommandReceiver launcherCommandReceiver; + @Autowired + TradeImporterService tradeImporterService; + @PostConstruct public void init() { super.init(); @@ -37,20 +45,20 @@ class TradeImporterServiceTest extends AbstractServiceTest { // Hazelcast.shutdownAll(); } + @BeforeEach + private void prepare(){ + clearAllInImdg(sTradesImdg); + } + /** - * {@link TradeImporterService#process()}
+ * {@link TradeImporterService#process(boolean)}
* Тест проверяет обновление сущности {@link STrades}.
* Входной запрос {@link LauncherCommandRequest}:
*/ -// @Test для работы теста нужна тестовая база Microsoft SQL с данными + @Test //для работы теста нужна тестовая база Microsoft SQL с данными void process() { - STrades trades = new STrades(); - trades.setId(22L); - trades.setTradeDate(LocalDate.of(2023,4,19)); - trades.setTradeNum(661486L); - trades.setOperation("operation20"); - trades.setClassCode("UESC"); - Long id = sTradesImdg.insert(trades); + STrades sTrade = getSTrade(); + Long id = sTradesImdg.insert(sTrade); 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)) .send(producerRecord.capture()); - STrades sTrades = sTradesImdg.getSingleObjectByFieldValues(Map.of("tradeDate", trades.getTradeDate(), - "tradeNum", trades.getTradeNum(), - "operation", trades.getOperation(), - "classCode", trades.getClassCode())); + STrades sTrades = tradeImporterService.getSTradesFromImdg(sTrade, sTradesImdg); assertEquals(sTrades.getId(), id); } + + /** + * {@link TradeImporterService#process(boolean)}
+ * Тест проверяет обновление сущности {@link STrades}.
+ * Входной запрос {@link LauncherCommandRequest}:
+ */ + @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; + } } \ No newline at end of file