trade-importer Оптимизация для большого количества S_TRADES

This commit is contained in:
AKurakin 2024-05-30 18:36:50 +03:00
parent f0cd944c64
commit 69e815f984
2 changed files with 35 additions and 13 deletions

View file

@ -49,6 +49,10 @@ public class STradesMapStore extends TemplateMapStore<STrades> {
return defaultLoadAllKeysOnTodayByField("trade_date", false);
}
public String[] getIndexingField() {
return new String[]{"tradeNum"};
}
@Override
public STrades objectReader(ResultSet resultSet) throws SQLException {
STrades object = new STrades();

View file

@ -29,6 +29,7 @@ import java.time.Instant;
import java.time.LocalDate;
import java.util.Collection;
import java.util.Map;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Consumer;
import java.util.function.Supplier;
@ -64,15 +65,30 @@ public class TradeImporterService {
}
public synchronized void process(boolean byCommand) {
log.debug("Start process import STrades from DB. byCommand={}", byCommand);
Collection<STrades> tradesFromDB = jdbcTemplate.query(String.format("SELECT * FROM %s.Trades", schema),
(resultSet, i) -> readSTrades(resultSet));
Instant currentInstant = Instant.now();
log.info("Start process import STrades from DB. byCommand={}", byCommand);
if (log.isDebugEnabled()) {
String countQuery = String.format("SELECT count(*) FROM %s.Trades", schema);
log.trace("Query: {}", countQuery);
Collection<Long> tradesCountFromDB = jdbcTemplate.query(countQuery,
(resultSet, i) -> resultSet.getObject(1, Long.class));
log.debug("Count of Trades in db: {}. Will be loading.", tradesCountFromDB);
}
log.debug("load from DB STrades: {}", tradesFromDB.size());
int created = 0;
int updated = 0;
for (STrades tradesDb : tradesFromDB) {
AtomicLong rows = new AtomicLong();
AtomicLong created = new AtomicLong();
AtomicLong updated = new AtomicLong();
final long logProgressTime = 60_000L; // интервал вывода в лог каждую минуту
AtomicLong logTime = new AtomicLong(System.currentTimeMillis() + logProgressTime);
String query = String.format("SELECT * FROM %s.Trades", schema);
log.debug("Load from DB STrades. Query: {}", query);
Instant currentInstant = Instant.now();
jdbcTemplate.query(query, resultSet -> {
rows.incrementAndGet();
if (System.currentTimeMillis() > logTime.get()) {
log.debug("Processed {} Trades rows...", rows.get());
logTime.set(System.currentTimeMillis() + logProgressTime);
}
STrades tradesDb = readSTrades(resultSet);
if (isValidTrades(tradesDb) && fillNessessaryFields(tradesDb)) {
STrades sTrades = getSTradesFromImdg(tradesDb, sTradesImdg);
if (sTrades != null && byCommand) {
@ -81,20 +97,22 @@ public class TradeImporterService {
tradesDb.setCreated(sTrades.getCreated());
tradesDb.setUpdated(currentInstant);
sTradesImdg.update(tradesDb);
updated++;
updated.incrementAndGet();
} else if (sTrades == null) {
// Создание объекта из БД
tradesDb.setCreated(currentInstant);
tradesDb.setUpdated(currentInstant);
sTradesImdg.insert(tradesDb);
created++;
created.incrementAndGet();
}
} else {
log.warn(messageResolver.resolve(new EnumMessage(sTradesNotValid, tradesDb.getTradeNum())));
}
}
log.debug("Created {} new STrades, {} updated.", created, updated);
if (created > 0 || updated > 0) {
});
log.debug("From DB {} STrades has loaded. Created {} new STrades, {} updated.",
rows.get(), created.get(), updated.get());
if (created.get() > 0 || updated.get() > 0) {
if (byCommand) {
log.debug("Successfully import STrades from DB by command. Send to kafka command, topic={}", S_TRADES_IMPORTED);
} else {