Поправил вывод в лог и чтение из БД.
This commit is contained in:
parent
bb396b71d5
commit
2f2835cc91
3 changed files with 71 additions and 7 deletions
|
|
@ -27,7 +27,6 @@ public class LauncherCommandReceiver extends QueueConsumer implements Initializi
|
||||||
callback(LauncherCommandRequest.class)
|
callback(LauncherCommandRequest.class)
|
||||||
.setConsumer(action -> importer.process(true))
|
.setConsumer(action -> importer.process(true))
|
||||||
.forDestination(Task.getOfTrades.topic(), callbacks::put); // GTRD
|
.forDestination(Task.getOfTrades.topic(), callbacks::put); // GTRD
|
||||||
importer.process(true);
|
|
||||||
init();
|
init();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -4,9 +4,7 @@ import org.apache.kafka.clients.producer.Producer;
|
||||||
import org.apache.kafka.clients.producer.ProducerRecord;
|
import org.apache.kafka.clients.producer.ProducerRecord;
|
||||||
import org.slf4j.Logger;
|
import org.slf4j.Logger;
|
||||||
import org.slf4j.LoggerFactory;
|
import org.slf4j.LoggerFactory;
|
||||||
import org.springframework.jdbc.core.BeanPropertyRowMapper;
|
|
||||||
import org.springframework.jdbc.core.JdbcTemplate;
|
import org.springframework.jdbc.core.JdbcTemplate;
|
||||||
import org.springframework.jdbc.core.RowMapper;
|
|
||||||
import org.springframework.scheduling.annotation.EnableScheduling;
|
import org.springframework.scheduling.annotation.EnableScheduling;
|
||||||
import org.springframework.scheduling.annotation.Scheduled;
|
import org.springframework.scheduling.annotation.Scheduled;
|
||||||
import org.springframework.stereotype.Service;
|
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.EnumMessage;
|
||||||
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
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.Collection;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
|
|
||||||
|
|
@ -33,7 +38,6 @@ public class TradeImporterService {
|
||||||
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 final IMessageResolver messageResolver;
|
||||||
private static final RowMapper<STrades> ROW_MAPPER = BeanPropertyRowMapper.newInstance(STrades.class);
|
|
||||||
|
|
||||||
|
|
||||||
public TradeImporterService(ImdgProvider imdgProvider, JdbcTemplate jdbcTemplate, Producer<String, Object> producer, IMessageResolver messageResolver) {
|
public TradeImporterService(ImdgProvider imdgProvider, JdbcTemplate jdbcTemplate, Producer<String, Object> producer, IMessageResolver messageResolver) {
|
||||||
|
|
@ -49,8 +53,10 @@ public class TradeImporterService {
|
||||||
}
|
}
|
||||||
|
|
||||||
public void process(boolean byCommand) {
|
public void process(boolean byCommand) {
|
||||||
Collection<STrades> tradesFromDB = jdbcTemplate.query("SELECT * FROM Trades", ROW_MAPPER);
|
log.debug("start process import STrades from DB");
|
||||||
|
Collection<STrades> 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) {
|
for (STrades tradesDb : tradesFromDB) {
|
||||||
if (isValidTrades(tradesDb)) {
|
if (isValidTrades(tradesDb)) {
|
||||||
STrades sTrades = getSTradesFromImdg(tradesDb, sTradesImdg);
|
STrades sTrades = getSTradesFromImdg(tradesDb, sTradesImdg);
|
||||||
|
|
@ -65,7 +71,7 @@ public class TradeImporterService {
|
||||||
STradesImportedRequest sTradesImportedRequest = new STradesImportedRequest();
|
STradesImportedRequest sTradesImportedRequest = new STradesImportedRequest();
|
||||||
producer.send(new ProducerRecord<>(S_TRADES_IMPORTED, sTradesImportedRequest));
|
producer.send(new ProducerRecord<>(S_TRADES_IMPORTED, sTradesImportedRequest));
|
||||||
log.debug("successfully import STrades from DB by command");
|
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<STrades> sTradesImdg) {
|
public STrades getSTradesFromImdg(STrades tradesDb, Imdg<STrades> sTradesImdg) {
|
||||||
|
|
@ -81,4 +87,63 @@ public class TradeImporterService {
|
||||||
&& StringUtils.hasText(trades.getOperation())
|
&& StringUtils.hasText(trades.getOperation())
|
||||||
&& trades.getTradeNum() != null;
|
&& 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;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -55,7 +55,7 @@ class TradeImporterServiceTest extends AbstractServiceTest {
|
||||||
* Тест проверяет обновление сущности {@link STrades}.<br>
|
* Тест проверяет обновление сущности {@link STrades}.<br>
|
||||||
* Входной запрос {@link LauncherCommandRequest}:<br>
|
* Входной запрос {@link LauncherCommandRequest}:<br>
|
||||||
*/
|
*/
|
||||||
@Test //для работы теста нужна тестовая база Microsoft SQL с данными
|
// @Test //для работы теста нужна тестовая база Microsoft SQL с данными
|
||||||
void process() {
|
void process() {
|
||||||
STrades sTrade = getSTrade();
|
STrades sTrade = getSTrade();
|
||||||
Long id = sTradesImdg.insert(sTrade);
|
Long id = sTradesImdg.insert(sTrade);
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue