diff --git a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/config/KafkaConfig.java b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/config/KafkaConfig.java index 0708fa158..79d7c78a5 100644 --- a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/config/KafkaConfig.java +++ b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/config/KafkaConfig.java @@ -1,15 +1,25 @@ package ru.spcex.clearing.trade.importer.config; import org.apache.kafka.clients.consumer.Consumer; -import org.apache.kafka.clients.producer.Producer; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.config.ConfigurableBeanFactory; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Scope; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.core.ProducerFactory; +import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; +import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings; +import ru.spcex.clearing.platform.messaging.service.RequestInfo; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.trade.importer.config.settings.ImportTradeServiceSettings; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgId; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.util.function.Supplier; @Configuration public class KafkaConfig { @@ -21,9 +31,30 @@ public class KafkaConfig { return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); } + @Bean + public ProducerFactory pf(ImportTradeServiceSettings settings) { + KafkaProducerSettings kafkaSettings = settings.getKafkaProducer(); + return KafkaProducerFactory.producerFactory(kafkaSettings); + } + + @Bean("kafkaTemplate") + public KafkaTemplate kafkaTemplate(ProducerFactory pf) { + return new KafkaTemplate<>(pf); + } + @Autowired @Bean - public Producer createProducer(ImportTradeServiceSettings settings) { - return KafkaProducerFactory.producer(settings.getKafkaProducer()); + public Supplier kafkaSender(KafkaTemplate kafkaTemplate, + ImdgProvider imdgProvider) { + ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); + return () -> KafkaSender + .setup() + .setKafkaTemplate(kafkaTemplate) + .idGenerator(imdgIdGenerator::nextId) + .imdgProvider(s -> { + Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); + return imdg::insert; + }) + .build(); } } 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 496b4b391..7b89a2c00 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 @@ -1,9 +1,8 @@ package ru.spcex.clearing.trade.importer.services; -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.beans.factory.annotation.Value; import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.scheduling.annotation.EnableScheduling; import org.springframework.scheduling.annotation.Scheduled; @@ -12,6 +11,7 @@ import org.springframework.util.StringUtils; import ru.clearing.classes.statics.data.misc.STrades; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.domain.cud.utilities.STradesImportedRequest; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; import ru.spcex.platform.utils.enumeration.EnumMessage; @@ -26,6 +26,7 @@ import java.time.Instant; import java.time.LocalDate; import java.util.Collection; import java.util.Map; +import java.util.function.Supplier; import static ru.spcex.clearing.platform.messaging.domain.Consts.S_TRADES_IMPORTED; import static ru.spcex.clearing.trade.importer.error.TradeImporterError.sTradesNotValid; @@ -36,14 +37,18 @@ public class TradeImporterService { private final Logger log = LoggerFactory.getLogger(getClass()); private final Imdg sTradesImdg; private final JdbcTemplate jdbcTemplate; - private final Producer producer; + private final Supplier kafka; private final IMessageResolver messageResolver; + @Value("${trade-importer.database.schema}") + private String schema; - public TradeImporterService(ImdgProvider imdgProvider, JdbcTemplate jdbcTemplate, Producer producer, IMessageResolver messageResolver) { + + public TradeImporterService(ImdgProvider imdgProvider, JdbcTemplate jdbcTemplate, Supplier kafka, IMessageResolver messageResolver) { + imdgProvider.waitAvailable(); this.sTradesImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_STrades, STrades.class); this.jdbcTemplate = jdbcTemplate; - this.producer = producer; + this.kafka = kafka; this.messageResolver = messageResolver; } @@ -54,7 +59,7 @@ public class TradeImporterService { public void process(boolean byCommand) { log.debug("Start process import STrades from DB"); - Collection tradesFromDB = jdbcTemplate.query("SELECT * FROM Trades", (resultSet, i) -> readSTrades(resultSet)); + Collection tradesFromDB = jdbcTemplate.query(String.format("SELECT * FROM %s.Trades", schema), (resultSet, i) -> readSTrades(resultSet)); Instant currentInstant = Instant.now(); log.debug("load from DB STrades: {}", tradesFromDB.size()); @@ -79,7 +84,7 @@ public class TradeImporterService { log.debug("Successfully import STrades from DB by scheduled. Send to kafka command, topic={}", S_TRADES_IMPORTED); } STradesImportedRequest sTradesImportedRequest = new STradesImportedRequest(); - producer.send(new ProducerRecord<>(S_TRADES_IMPORTED, sTradesImportedRequest)); + kafka.get().sendRequestToQueue(S_TRADES_IMPORTED, sTradesImportedRequest); } public STrades getSTradesFromImdg(STrades tradesDb, Imdg sTradesImdg) { diff --git a/clearing-parent/trade-importer/src/main/resources/application.properties b/clearing-parent/trade-importer/src/main/resources/application.properties index e2cfdbd94..7b7f8f668 100644 --- a/clearing-parent/trade-importer/src/main/resources/application.properties +++ b/clearing-parent/trade-importer/src/main/resources/application.properties @@ -5,7 +5,8 @@ trade-importer.cron.load-from-db-cron=0 0/5 * * * ? trade-importer.database.login=sa trade-importer.database.password=Aa123456 -trade-importer.database.url=jdbc:sqlserver://localhost:1433;database=SPVB_TS;schema=dbo +trade-importer.database.schema=SPVB_TS +trade-importer.database.url=jdbc:sqlserver://10.200.200.144:1433;database=ni; trade-importer.database.driver=com.microsoft.sqlserver.jdbc.SQLServerDriver trade-importer.hazelcast.cluster-members=127.0.0.1:5701