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 ade26b866..d66ed5f21 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 @@ -59,11 +59,11 @@ public class TradeImporterService { } @Scheduled(cron = "${trade-importer.cron.load-from-db-cron}") - public void run() { + public synchronized void run() { process(false); } - public void process(boolean byCommand) { + public synchronized void process(boolean byCommand) { log.debug("Start process import STrades from DB"); Collection tradesFromDB = jdbcTemplate.query(String.format("SELECT * FROM %s.Trades", schema), (resultSet, i) -> readSTrades(resultSet)); @@ -75,7 +75,7 @@ public class TradeImporterService { for (STrades tradesDb : tradesFromDB) { if (isValidTrades(tradesDb) && fillNessessaryFields(tradesDb)) { STrades sTrades = getSTradesFromImdg(tradesDb, sTradesImdg); - if (sTrades != null) { + if (sTrades != null && byCommand) { // Обновление всех полей объекта из БД tradesDb.setId(sTrades.getId()); tradesDb.setCreated(sTrades.getCreated()); @@ -94,14 +94,17 @@ public class TradeImporterService { } } log.debug("Created {} new STrades, {} updated.", created, updated); - - if (byCommand) { - log.debug("Successfully import STrades from DB by command. Send to kafka command, topic={}", S_TRADES_IMPORTED); + if (created > 0 || updated > 0) { + if (byCommand) { + log.debug("Successfully import STrades from DB by command. Send to kafka command, topic={}", S_TRADES_IMPORTED); + } else { + log.debug("Successfully import STrades from DB by scheduled. Send to kafka command, topic={}", S_TRADES_IMPORTED); + } + STradesImportedRequest sTradesImportedRequest = new STradesImportedRequest(); + kafka.get().sendRequestToQueue(S_TRADES_IMPORTED, sTradesImportedRequest); } else { - log.debug("Successfully import STrades from DB by scheduled. Send to kafka command, topic={}", S_TRADES_IMPORTED); + log.debug("No STrades were created or updated from DB. Kafka command will not be send"); } - STradesImportedRequest sTradesImportedRequest = new STradesImportedRequest(); - kafka.get().sendRequestToQueue(S_TRADES_IMPORTED, sTradesImportedRequest); } /**