Поправил настройки подключения к БД и отправку в кафку.
This commit is contained in:
parent
0dd8b6840c
commit
4bff19f9ba
3 changed files with 48 additions and 11 deletions
|
|
@ -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<String, Object> pf(ImportTradeServiceSettings settings) {
|
||||
KafkaProducerSettings kafkaSettings = settings.getKafkaProducer();
|
||||
return KafkaProducerFactory.producerFactory(kafkaSettings);
|
||||
}
|
||||
|
||||
@Bean("kafkaTemplate")
|
||||
public KafkaTemplate<String, Object> kafkaTemplate(ProducerFactory<String, Object> pf) {
|
||||
return new KafkaTemplate<>(pf);
|
||||
}
|
||||
|
||||
@Autowired
|
||||
@Bean
|
||||
public Producer<String, Object> createProducer(ImportTradeServiceSettings settings) {
|
||||
return KafkaProducerFactory.producer(settings.getKafkaProducer());
|
||||
public Supplier<KafkaSender> kafkaSender(KafkaTemplate<String, Object> kafkaTemplate,
|
||||
ImdgProvider imdgProvider) {
|
||||
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
|
||||
return () -> KafkaSender
|
||||
.setup()
|
||||
.setKafkaTemplate(kafkaTemplate)
|
||||
.idGenerator(imdgIdGenerator::nextId)
|
||||
.imdgProvider(s -> {
|
||||
Imdg<RequestInfo> imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);
|
||||
return imdg::insert;
|
||||
})
|
||||
.build();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<STrades> sTradesImdg;
|
||||
private final JdbcTemplate jdbcTemplate;
|
||||
private final Producer<String, Object> producer;
|
||||
private final Supplier<KafkaSender> kafka;
|
||||
private final IMessageResolver messageResolver;
|
||||
|
||||
@Value("${trade-importer.database.schema}")
|
||||
private String schema;
|
||||
|
||||
public TradeImporterService(ImdgProvider imdgProvider, JdbcTemplate jdbcTemplate, Producer<String, Object> producer, IMessageResolver messageResolver) {
|
||||
|
||||
public TradeImporterService(ImdgProvider imdgProvider, JdbcTemplate jdbcTemplate, Supplier<KafkaSender> 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<STrades> tradesFromDB = jdbcTemplate.query("SELECT * FROM Trades", (resultSet, i) -> readSTrades(resultSet));
|
||||
Collection<STrades> 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<STrades> sTradesImdg) {
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue