ialbert 2022-09-22 19:56:11 +03:00
parent 29c818a76a
commit d147732b0e
7 changed files with 192 additions and 7 deletions

View file

@ -1,6 +1,7 @@
package ru.spcex.clearing.balance.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;
@ -8,13 +9,20 @@ import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Scope;
import ru.spcex.clearing.balance.config.element.BalanceServiceSettings;
import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory;
import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory;
@Configuration
public class KafkaConfig {
@Autowired
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
@Bean
public Consumer<String, Object> createProducer(BalanceServiceSettings settings) {
return KafkaConsumerFactory.consumer(settings.getKafka());
public Consumer<String, Object> createConsumer(BalanceServiceSettings settings) {
return KafkaConsumerFactory.consumer(settings.getKafkaConsumer());
}
@Autowired
@Bean
public Producer<String, Object> createProducer(BalanceServiceSettings settings) {
return KafkaProducerFactory.producer(settings.getKafkaProducer());
}
}

View file

@ -0,0 +1,30 @@
package ru.spcex.clearing.balance.config;
import org.apache.kafka.clients.producer.Producer;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
@Configuration
public class KafkaSenderConfig {
@Autowired
@Bean
public KafkaSender kafkaSender(Producer<String, Object> kafkaProducer, ImdgProvider imdgProvider) {
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
return KafkaSender
.setup()
.producer(kafkaProducer)
.idGenerator(imdgIdGenerator::nextId)
.imdgProvider(s -> {
Imdg<RequestInfo> imdg = imdgProvider.getImdg(s, RequestInfo.class);
return imdg::insert;
})
.build();
}
}

View file

@ -4,6 +4,7 @@ import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.PropertySource;
import org.springframework.stereotype.Component;
import ru.spcex.clearing.platform.messaging.config.element.KafkaConsumerSettings;
import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings;
import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams;
@Component
@ -11,7 +12,8 @@ import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams;
@ConfigurationProperties("balance-service")
public class BalanceServiceSettings {
private HazelcastClientParams hazelcast;
private KafkaConsumerSettings kafka;
private KafkaConsumerSettings kafkaConsumer;
private KafkaProducerSettings kafkaProducer;
public HazelcastClientParams getHazelcast() {
return hazelcast;
@ -21,11 +23,19 @@ public class BalanceServiceSettings {
this.hazelcast = hazelcast;
}
public KafkaConsumerSettings getKafka() {
return kafka;
public KafkaConsumerSettings getKafkaConsumer() {
return kafkaConsumer;
}
public void setKafka(KafkaConsumerSettings kafka) {
this.kafka = kafka;
public void setKafkaConsumer(KafkaConsumerSettings kafkaConsumer) {
this.kafkaConsumer = kafkaConsumer;
}
public KafkaProducerSettings getKafkaProducer() {
return kafkaProducer;
}
public void setKafkaProducer(KafkaProducerSettings kafkaProducer) {
this.kafkaProducer = kafkaProducer;
}
}

View file

@ -0,0 +1,7 @@
package ru.spcex.clearing.platform.messaging.service.sender;
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
public interface KafkaImdgInsert {
void insert(RequestInfo requestInfo);
}

View file

@ -0,0 +1,76 @@
package ru.spcex.clearing.platform.messaging.service.sender;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import ru.spcex.clearing.platform.messaging.domain.ActionType;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
import ru.spcex.platform.classes.base.SpcexObjectBase;
import ru.spcex.platform.utils.log.ExceptionUtils;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Future;
import java.util.function.Function;
import java.util.function.Supplier;
public class KafkaSender {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Map<String, KafkaImdgInsert> allImdgMaps;
private Producer<String, Object> kafka;
private Supplier<Long> idGenerator;
private Function<String, KafkaImdgInsert> imdgGenerator;
KafkaSender() {
this.allImdgMaps = new ConcurrentHashMap<>();
}
public static KafkaSenderBuilderImpl setup() {
return new KafkaSenderBuilderImpl();
}
public Long sendRequestToQueue(String destination, Object requestPayload) {
BaseRequest<Object> request = new BaseRequest<>();
request.setId(idGenerator.get());
request.setActionType(ActionType.SYSTEM);
request.setRequestPayload(requestPayload);
//сохраняет данные о запросе в хранилище
saveRequestToStorage(destination, request);
Future<RecordMetadata> send = kafka.send(new ProducerRecord<>(destination, request));
try {
send.get();
} catch (InterruptedException | ExecutionException e) {
log.error(ExceptionUtils.getStackTrace(e));
return null;
}
return request.getId();
}
private void saveRequestToStorage(String destination, BaseRequest<Object> request) {
KafkaImdgInsert imdgInsert = getImdg(destination);
RequestInfo requestInfo = RequestInfo.create(request.getId());
imdgInsert.insert(requestInfo);
}
@SuppressWarnings("unchecked")
private <T extends SpcexObjectBase> KafkaImdgInsert getImdg(String mapName) {
return allImdgMaps.computeIfAbsent(mapName, (mapName1) -> imdgGenerator.apply(mapName));
}
void setProducer(Producer<String, Object> kafka) {
this.kafka = kafka;
}
void setIdGenerator(Supplier<Long> idGenerator) {
this.idGenerator = idGenerator;
}
void setImdgProvider(Function<String, KafkaImdgInsert> imdgGenerator) {
this.imdgGenerator = imdgGenerator;
}
}

View file

@ -0,0 +1,16 @@
package ru.spcex.clearing.platform.messaging.service.sender;
import org.apache.kafka.clients.producer.Producer;
import java.util.function.Function;
import java.util.function.Supplier;
public interface KafkaSenderBuilder {
KafkaSenderBuilder producer(Producer<String, Object> kafka);
KafkaSenderBuilder idGenerator(Supplier<Long> idGenerator);
KafkaSenderBuilder imdgProvider(Function<String, KafkaImdgInsert> imdgGenerator);
KafkaSender build();
}

View file

@ -0,0 +1,38 @@
package ru.spcex.clearing.platform.messaging.service.sender;
import org.apache.kafka.clients.producer.Producer;
import java.util.function.Function;
import java.util.function.Supplier;
public class KafkaSenderBuilderImpl implements KafkaSenderBuilder{
private final KafkaSender kafkaSender;
public KafkaSenderBuilderImpl() {
this.kafkaSender = new KafkaSender();
}
@Override
public KafkaSenderBuilder producer(Producer<String, Object> kafka) {
kafkaSender.setProducer(kafka);
return this;
}
@Override
public KafkaSenderBuilder idGenerator(Supplier<Long> idGenerator) {
kafkaSender.setIdGenerator(idGenerator);
return this;
}
@Override
public KafkaSenderBuilder imdgProvider(Function<String, KafkaImdgInsert> imdgGenerator) {
kafkaSender.setImdgProvider(imdgGenerator);
return this;
}
@Override
public KafkaSender build() {
return kafkaSender;
}
}