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 5c37e06be..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 @@ -27,7 +27,6 @@ public class LauncherCommandReceiver extends QueueConsumer implements Initializi callback(LauncherCommandRequest.class) .setConsumer(action -> importer.process(true)) .forDestination(Task.getOfTrades.topic(), callbacks::put); // GTRD - importer.process(true); 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 12ff78487..84dd75c82 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 @@ -4,9 +4,7 @@ import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.ProducerRecord; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.springframework.jdbc.core.BeanPropertyRowMapper; import org.springframework.jdbc.core.JdbcTemplate; -import org.springframework.jdbc.core.RowMapper; import org.springframework.scheduling.annotation.EnableScheduling; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; @@ -19,6 +17,13 @@ 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.Map; @@ -33,7 +38,6 @@ public class TradeImporterService { 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, IMessageResolver messageResolver) { @@ -49,8 +53,10 @@ public class TradeImporterService { } public void process(boolean byCommand) { - Collection tradesFromDB = jdbcTemplate.query("SELECT * FROM Trades", ROW_MAPPER); + log.debug("start process import STrades from DB"); + Collection tradesFromDB = jdbcTemplate.query("SELECT * FROM Trades", (resultSet, i) -> readSTrades(resultSet)); + log.debug(String.format("load from DB STrades: %d", tradesFromDB.size())); for (STrades tradesDb : tradesFromDB) { if (isValidTrades(tradesDb)) { STrades sTrades = getSTradesFromImdg(tradesDb, sTradesImdg); @@ -65,7 +71,7 @@ public class TradeImporterService { STradesImportedRequest sTradesImportedRequest = new STradesImportedRequest(); producer.send(new ProducerRecord<>(S_TRADES_IMPORTED, sTradesImportedRequest)); log.debug("successfully import STrades from DB by command"); - } + }else log.debug("successfully import STrades from DB by scheduled"); } public STrades getSTradesFromImdg(STrades tradesDb, Imdg sTradesImdg) { @@ -81,4 +87,63 @@ public class TradeImporterService { && StringUtils.hasText(trades.getOperation()) && trades.getTradeNum() != null; } + + public STrades readSTrades(ResultSet resultSet) throws SQLException { + STrades sTrades = new STrades(); + sTrades.setTradeNum(resultSet.getObject("TradeNum", Long.class)); + sTrades.setOperation(resultSet.getObject("Operation", String.class)); + sTrades.setClassCode(resultSet.getObject("ClassCode", String.class)); + sTrades.setTradeDate(getLocalDateFromSqlDate(resultSet, "TradeDate")); + sTrades.setSecCode(resultSet.getObject("SecCode", String.class)); + sTrades.setAccruedint(resultSet.getObject("Accruedint", BigDecimal.class)); + sTrades.setAccruedint2(resultSet.getObject("Accruedint2", BigDecimal.class)); + sTrades.setLowerDiscount(resultSet.getObject("LowerDiscount", BigDecimal.class)); + sTrades.setOrderNum(resultSet.getObject("OrderNum", Long.class)); + sTrades.setPrice(resultSet.getObject("Price", BigDecimal.class)); + sTrades.setPrice2(resultSet.getObject("Price2", BigDecimal.class)); + sTrades.setRepoRate(resultSet.getObject("RepoRate", BigDecimal.class)); + sTrades.setRepoValue(resultSet.getObject("RepoValue", BigDecimal.class)); + sTrades.setRepo2Value(resultSet.getObject("Repo2Value", BigDecimal.class)); + sTrades.setStartDiscount(resultSet.getObject("StartDiscount", BigDecimal.class)); + sTrades.setTsCommission(resultSet.getObject("TSCommission", BigDecimal.class)); + sTrades.setUpperDiscount(resultSet.getObject("UpperDiscount", BigDecimal.class)); + sTrades.setValue(resultSet.getObject("Value", BigDecimal.class)); + sTrades.setYield(resultSet.getObject("Yield", BigDecimal.class)); + sTrades.setQty(resultSet.getObject("Qty", BigDecimal.class)); + sTrades.setQtyPcs(resultSet.getObject("Qty_pcs", BigDecimal.class)); + sTrades.setTradeDateTime(getInstantFromTimestamp(resultSet, "TradeDateTime")); + sTrades.setRepoTerm(resultSet.getObject("RepoTerm", Long.class)); + sTrades.setClearingCommission(resultSet.getObject("ClearingCommission", BigDecimal.class)); + sTrades.setExchangeCommission(resultSet.getObject("ExchangeCommission", BigDecimal.class)); + sTrades.setTechCenterCommission(resultSet.getObject("TechCenterCommission", BigDecimal.class)); + sTrades.setAccount(resultSet.getObject("Account", String.class)); + sTrades.setBrokerRef(resultSet.getObject("BrokerRef", String.class)); + sTrades.setClientCode(resultSet.getObject("ClientCode", String.class)); + sTrades.setSettleCode(resultSet.getObject("SettleCode", String.class)); + sTrades.setUserId(resultSet.getObject("UserId", String.class)); + sTrades.setExchangeCode(resultSet.getObject("ExchangeCode", String.class)); + sTrades.setFirmId(resultSet.getObject("FirmId", String.class)); + sTrades.setFirmName(resultSet.getObject("FirmName", String.class)); + sTrades.setCpFirmId(resultSet.getObject("CPFirmId", String.class)); + sTrades.setCpFirmName(resultSet.getObject("CPFirmName", String.class)); + sTrades.setClassName(resultSet.getObject("ClassName", String.class)); + sTrades.setSecName(resultSet.getObject("SecName", String.class)); + sTrades.setSettleDate(getLocalDateFromSqlDate(resultSet, "SettleDate")); + sTrades.setSettleCurrency(resultSet.getObject("SettleCurrency", String.class)); + sTrades.setTradeCurrency(resultSet.getObject("TradeCurrency", String.class)); + sTrades.setTradeTimeMs(resultSet.getObject("TradeTimeMs", Long.class)); + sTrades.setBankAccId(resultSet.getObject("BankAccId", String.class)); + sTrades.setSection(resultSet.getObject("Section", String.class)); + return sTrades; + } + + 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; + } } 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 49068bcb9..fdd014270 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 @@ -55,7 +55,7 @@ class TradeImporterServiceTest extends AbstractServiceTest { * Тест проверяет обновление сущности {@link STrades}.
* Входной запрос {@link LauncherCommandRequest}:
*/ - @Test //для работы теста нужна тестовая база Microsoft SQL с данными +// @Test //для работы теста нужна тестовая база Microsoft SQL с данными void process() { STrades sTrade = getSTrade(); Long id = sTradesImdg.insert(sTrade);