From 69e815f98460ce473a9b4bd78bba589bfa5e60c9 Mon Sep 17 00:00:00 2001 From: AKurakin Date: Thu, 30 May 2024 18:36:50 +0300 Subject: [PATCH] =?UTF-8?q?trade-importer=20=D0=9E=D0=BF=D1=82=D0=B8=D0=BC?= =?UTF-8?q?=D0=B8=D0=B7=D0=B0=D1=86=D0=B8=D1=8F=20=D0=B4=D0=BB=D1=8F=20?= =?UTF-8?q?=D0=B1=D0=BE=D0=BB=D1=8C=D1=88=D0=BE=D0=B3=D0=BE=20=D0=BA=D0=BE?= =?UTF-8?q?=D0=BB=D0=B8=D1=87=D0=B5=D1=81=D1=82=D0=B2=D0=B0=20S=5FTRADES?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../clearing/imdg/object/STradesMapStore.java | 4 ++ .../services/TradeImporterService.java | 44 +++++++++++++------ 2 files changed, 35 insertions(+), 13 deletions(-) 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 6a8c95582..ec3645a21 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 @@ -49,6 +49,10 @@ public class STradesMapStore extends TemplateMapStore { return defaultLoadAllKeysOnTodayByField("trade_date", false); } + public String[] getIndexingField() { + return new String[]{"tradeNum"}; + } + @Override public STrades objectReader(ResultSet resultSet) throws SQLException { STrades object = new STrades(); 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 5454cee89..1141631b8 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 @@ -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 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 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 {