http://jira.mfd.msk:8088/browse/CLS-290 добавил Spring KafkaTemplate

This commit is contained in:
ialbert 2023-05-15 12:04:05 +03:00
parent ce84c28297
commit 6c23cbca36
22 changed files with 478 additions and 70 deletions

View file

@ -7,10 +7,13 @@ 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.account.config.settings.AccountServiceSettings;
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.platform.imdg.api.Imdg;
@ -33,13 +36,24 @@ public class KafkaConfig {
return KafkaProducerFactory.producer(settings.getKafkaProducer());
}
@Bean
public ProducerFactory<String, Object> pf(AccountServiceSettings settings) {
KafkaProducerSettings kafkaSettings = settings.getKafkaProducer();
return KafkaProducerFactory.producerFactory(kafkaSettings);
}
@Bean
public KafkaTemplate<String, Object> kafkaTemplate(ProducerFactory<String, Object> pf) {
return new KafkaTemplate<>(pf);
}
@Autowired
@Bean
public KafkaSender kafkaSender(Producer<String, Object> kafkaProducer, ImdgProvider imdgProvider) {
public KafkaSender kafkaSender(KafkaTemplate<String, Object> kafkaTemplate, ImdgProvider imdgProvider) {
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
return KafkaSender
.setup()
.producer(kafkaProducer)
.setKafkaTemplate(kafkaTemplate)
.idGenerator(imdgIdGenerator::nextId)
.imdgProvider(s -> {
Imdg<RequestInfo> imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);

View file

@ -8,10 +8,13 @@ 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.balance.config.element.BalanceServiceSettings;
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.platform.imdg.api.Imdg;
@ -33,14 +36,25 @@ public class KafkaConfig {
return KafkaProducerFactory.producer(settings.getKafkaProducer());
}
@Bean
public ProducerFactory<String, Object> pf(BalanceServiceSettings 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 KafkaSender kafkaSender(@Qualifier("kafkaProducer") Producer<String, Object> kafkaProducer,
public KafkaSender kafkaSender(@Qualifier("kafkaTemplate") KafkaTemplate<String, Object> kafkaTemplate,
ImdgProvider imdgProvider) {
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
return KafkaSender
.setup()
.producer(kafkaProducer)
.setKafkaTemplate(kafkaTemplate)
.idGenerator(imdgIdGenerator::nextId)
.imdgProvider(s -> {
Imdg<RequestInfo> imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);

View file

@ -7,10 +7,13 @@ 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.config.element.ClearingServiceSettings;
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.platform.imdg.api.Imdg;
@ -32,13 +35,24 @@ public class KafkaConfig {
return KafkaProducerFactory.producer(settings.getKafkaProducer());
}
@Bean
public ProducerFactory<String, Object> pf(ClearingServiceSettings 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 KafkaSender kafkaSender(Producer<String, Object> kafkaProducer, ImdgProvider imdgProvider) {
public KafkaSender kafkaSender(KafkaTemplate<String, Object> kafkaTemplate, ImdgProvider imdgProvider) {
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
return KafkaSender
.setup()
.producer(kafkaProducer)
.setKafkaTemplate(kafkaTemplate)
.idGenerator(imdgIdGenerator::nextId)
.imdgProvider(s -> {
Imdg<RequestInfo> imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);
@ -49,11 +63,11 @@ public class KafkaConfig {
@Autowired
@Bean("kafkaSenderWithoutRequestInfo")
public KafkaSender kafkaSenderWithoutRequestInfo(Producer<String, Object> kafkaProducer, ImdgProvider imdgProvider) {
public KafkaSender kafkaSenderWithoutRequestInfo(KafkaTemplate<String, Object> kafkaTemplate, ImdgProvider imdgProvider) {
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
return KafkaSender
.setup()
.producer(kafkaProducer)
.setKafkaTemplate(kafkaTemplate)
.idGenerator(imdgIdGenerator::nextId)
.saveRequestInfo(false)
.build();

View file

@ -3,6 +3,7 @@ package ru.spcex.clearing.service;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.company.ClearingMemberCategory;
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
@ -62,7 +63,7 @@ public class Clearing {
@Autowired
public Clearing(ImdgProvider imdgProvider,
ExecutionDepositSorter sorter, PaymentInstructionCreator paymentInstructionCreator,
BiFunction<ClearingCategory, ExecutionDeposit, IValidator> validation, KafkaSender kafka, LiabilitiesClaimsAssetsCreator liabilitiesClaimsAssetsCreator, LiabilitiesClaimsMoneyCreator lbltsClmsMoneyCreator) {
BiFunction<ClearingCategory, ExecutionDeposit, IValidator> validation, @Qualifier("kafkaSender") KafkaSender kafka, LiabilitiesClaimsAssetsCreator liabilitiesClaimsAssetsCreator, LiabilitiesClaimsMoneyCreator lbltsClmsMoneyCreator) {
this.executionDepositImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionDeposit, ExecutionDeposit.class);
this.clearingCategoryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingMemberCategory, ClearingMemberCategory.class);
this.liabilitiesClaimsAssetsImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_LiabilitiesClaimsAssets, LiabilitiesClaimsAssets.class);

View file

@ -7,9 +7,18 @@ 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.company.config.settings.CompanyServiceSettings;
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.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
@Configuration
public class KafkaConfig {
@ -26,4 +35,31 @@ public class KafkaConfig {
return KafkaProducerFactory.producer(settings.getKafkaProducer());
}
@Bean
public ProducerFactory<String, Object> pf(CompanyServiceSettings 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 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();
}
}

View file

@ -1,32 +0,0 @@
package ru.spcex.clearing.company.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.imdg.IMDGDistributedNames;
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(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);
return imdg::insert;
})
.build();
}
}

View file

@ -1,12 +1,16 @@
package ru.spcex.clearing.dbf.exporter.config;
import org.apache.kafka.clients.producer.Producer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.ProducerFactory;
import ru.spcex.clearing.dbf.exporter.config.settings.ExportDBFServiceSettings;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
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.platform.imdg.api.Imdg;
@ -19,28 +23,42 @@ public class KafkaSenderConfig {
Logger log = LoggerFactory.getLogger(getClass());
private final ImdgProvider imdgProvider;
private Producer<String, Object> kafkaProducer;
@Autowired
public KafkaSenderConfig(ImdgProvider imdgProvider) {
this.imdgProvider = imdgProvider;
}
@Autowired(required = false)
public void setKafkaProducer(Producer<String, Object> kafkaProducer) {
this.kafkaProducer = kafkaProducer;
@Bean
public ProducerFactory<String, Object> pf(ExportDBFServiceSettings settings) {
if (settings.getKafkaProducer() == null) {
return null;
}
KafkaProducerSettings kafkaSettings = settings.getKafkaProducer();
return KafkaProducerFactory.producerFactory(kafkaSettings);
}
@Autowired(required = false)
@Bean("kafkaTemplate")
public KafkaTemplate<String, Object> kafkaTemplate(ProducerFactory<String, Object> pf) {
if (pf == null) {
return null;
}
return new KafkaTemplate<>(pf);
}
@Autowired(required = false)
@Bean
public KafkaSender kafkaSender() {
if (kafkaProducer == null || imdgProvider == null) {
log.info("Can not create KafkaSender: kafkaProducer={}, imdgProvider={}", kafkaProducer, imdgProvider);
public KafkaSender kafkaSender(KafkaTemplate<String, Object> kafkaTemplate,
ImdgProvider imdgProvider) {
if (kafkaTemplate == null) {
log.info("Can not create KafkaSender: no kafka-producer settings");
return null;
}
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
return KafkaSender
.setup()
.producer(kafkaProducer)
.setKafkaTemplate(kafkaTemplate)
.idGenerator(imdgIdGenerator::nextId)
.imdgProvider(s -> {
Imdg<RequestInfo> imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);

View file

@ -1,16 +1,18 @@
package ru.spcex.clearing.dbf.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.dbf.importer.config.settings.ImportDBFServiceSettings;
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.platform.imdg.api.Imdg;
@ -22,22 +24,31 @@ import java.util.function.Supplier;
@Configuration
public class KafkaConfig {
@Bean
public ProducerFactory<String, Object> pf(ImportDBFServiceSettings 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 Supplier<KafkaSender> kafkaSender(ImportDBFServiceSettings settings, ImdgProvider imdgProvider) {
return () -> {
Producer<String, Object> kafkaProducer = KafkaProducerFactory.producer(settings.getKafkaProducer());
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
return KafkaSender
.setup()
.producer(kafkaProducer)
.idGenerator(imdgIdGenerator::nextId)
.imdgProvider(s -> {
Imdg<RequestInfo> imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);
return imdg::insert;
})
.build();
};
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();
}
@Autowired

View file

@ -1,9 +1,11 @@
package ru.spcex.platform.imdg.iml.hazelcast.adapter;
import com.hazelcast.aggregation.Aggregators;
import com.hazelcast.core.HazelcastInstance;
import com.hazelcast.core.ILock;
import com.hazelcast.core.IMap;
import com.hazelcast.core.IdGenerator;
import com.hazelcast.projection.Projections;
import com.hazelcast.query.Predicate;
import com.hazelcast.query.Predicates;
import com.hazelcast.query.SqlPredicate;
@ -34,6 +36,35 @@ public class ImdgHazelcast<T extends SpcexObjectBase> implements Imdg<T> {
return map.values();
}
@Override
public Long aggregateLongMax(String fieldName, ImdgPredicate predicate) {
Predicate<Long, T> p = ((ImdgPredicateHazelcast) predicate).getRawPredicate();
return map.aggregate(Aggregators.longMax(fieldName), p);
}
@Override
public T aggregateByMax(String fieldName, ImdgPredicate predicate) {
Map.Entry<Long, T> entry;
if (predicate != null) {
Predicate<Long, T> p = ((ImdgPredicateHazelcast) predicate).getRawPredicate();
entry = map.aggregate(Aggregators.maxBy(fieldName), p);
} else {
entry = map.aggregate(Aggregators.maxBy(fieldName));
}
return entry != null ? entry.getValue() : null;
}
@Override
public <A> Collection<A> projectSingleAttribute(String fieldName, ImdgPredicate predicate) {
if (predicate != null) {
Predicate<Long, T> p = null;
p = ((ImdgPredicateHazelcast) predicate).getRawPredicate();
return map.project(Projections.singleAttribute(fieldName), p);
} else {
return map.project(Projections.singleAttribute(fieldName));
}
}
@Override
public Integer size() {
return map.size();

View file

@ -25,6 +25,10 @@
<groupId>ru.spcex.platform</groupId>
<artifactId>platform-classes-base</artifactId>
</dependency>
<dependency>
<groupId>ru.spcex.platform</groupId>
<artifactId>platform-enum</artifactId>
</dependency>
<dependency>
<groupId>ru.spcex.platform</groupId>
<artifactId>platform-utils</artifactId>

View file

@ -64,6 +64,34 @@ public interface Imdg<T extends SpcexObjectBase> {
throw new UnsupportedOperationException("not implemented projectionsAttributeBySql");
}
/**
* при необходимости ввести ImdgAggregator
*/
default Long aggregateLongMax(String fieldName, ImdgPredicate predicate) {
throw new UnsupportedOperationException("not implemented aggregate");
}
/**
* при необходимости ввести ImdgAggregator
*/
default T aggregateByMax(String fieldName, ImdgPredicate predicate) {
throw new UnsupportedOperationException("not implemented aggregate");
}
/**
* при необходимости ввести ImdgProjection
*/
default <A> Collection<A> projectSingleAttribute(String fieldName, ImdgPredicate predicate) {
throw new UnsupportedOperationException("not implemented projectSingleAttribute");
}
/**
* при необходимости ввести ImdgProjection
*/
default <A> Collection<A> projectSingleAttribute(String fieldName) {
return projectSingleAttribute(fieldName, null);
}
default Integer size() {
throw new UnsupportedOperationException("not implemented size");
}

View file

@ -0,0 +1,18 @@
package ru.spcex.platform.imdg.api.predicate.specific;
import ru.spcex.platform.enumeration.StatementType;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
import java.util.Collection;
public class StatementRevisePredicate {
public static ImdgPredicate get(Imdg<?> pb, Collection<Long> currenciesId) {
ImdgPredicateBuilder prdctBuilder = pb.predicateBuilder();
return prdctBuilder.and(
prdctBuilder.in("currencyId", currenciesId.toArray(new Comparable[0])),
prdctBuilder.equals("statementType", StatementType.incr.getKey())
);
}
}

View file

@ -2,8 +2,14 @@ package ru.spcex.clearing.platform.messaging.config;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import org.springframework.kafka.support.serializer.ErrorHandlingDeserializer;
import ru.spcex.clearing.platform.messaging.config.element.KafkaConsumerSettings;
import ru.spcex.clearing.platform.messaging.serialization.JsonDeserializer;
import java.util.HashMap;
import java.util.Map;
import java.util.Properties;
public class KafkaConsumerFactory {
@ -20,4 +26,24 @@ public class KafkaConsumerFactory {
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
return new KafkaConsumer<>(props);
}
public static ConsumerFactory<String, Object> consumerFactory(KafkaConsumerSettings kafkaSettings) {
Map<String, Object> props = new HashMap<>();
//Properties props = new Properties();
props.put("bootstrap.servers", kafkaSettings.getBootstrapServers());
if (kafkaSettings.getGroupId() != null && kafkaSettings.getGroupId().length() > 0) {
props.put("group.id", kafkaSettings.getGroupId());
}
props.put("enable.auto.commit", kafkaSettings.getEnableAutoCommit().toString());
props.put("session.timeout.ms", kafkaSettings.getSessionTimeoutMs().toString());
props.put("auto.offset.reset", kafkaSettings.getAutoOffsetReset());
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
// props.put("value.deserializer", JsonDeserializer.class.getName());
props.put("value.deserializer", ErrorHandlingDeserializer.class.getName());
props.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, JsonDeserializer.class.getName());
DefaultKafkaConsumerFactory<String, Object> cf = new DefaultKafkaConsumerFactory<>(props);
cf.setValueDeserializer(new ErrorHandlingDeserializer<>(new JsonDeserializer()));
return cf;
}
}

View file

@ -3,9 +3,13 @@ package ru.spcex.clearing.platform.messaging.config;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
import org.springframework.kafka.core.ProducerFactory;
import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings;
import ru.spcex.clearing.platform.messaging.serialization.JsonSerializer;
import java.util.HashMap;
import java.util.Map;
import java.util.Properties;
public class KafkaProducerFactory {
@ -37,4 +41,29 @@ public class KafkaProducerFactory {
return new KafkaProducer<>(kafkaProps);
}
public static ProducerFactory<String, Object> producerFactory(KafkaProducerSettings kafkaSettings) {
Map<String, Object> kafkaProps = new HashMap<>();
////Assign localhost id
kafkaProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaSettings.getBootstrapServers());
//
////Set acknowledgements for producer requests.
kafkaProps.put(ProducerConfig.ACKS_CONFIG, kafkaSettings.getAcks());
//
////If the request fails, the producer can automatically retry,
kafkaProps.put(ProducerConfig.RETRIES_CONFIG, kafkaSettings.getRetries());
kafkaProps.put(ProducerConfig.RECONNECT_BACKOFF_MS_CONFIG, 10000);
//
////Specify buffer size in config
kafkaProps.put(ProducerConfig.BATCH_SIZE_CONFIG, kafkaSettings.getBatchSize());
//
////Reduce the no of requests less than 0
kafkaProps.put(ProducerConfig.LINGER_MS_CONFIG, kafkaSettings.getLingerMs());
////The buffer.memory controls the total amount of memory available to the producer for buffering.
kafkaProps.put(ProducerConfig.BUFFER_MEMORY_CONFIG, kafkaSettings.getBufferMemory());
kafkaProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
kafkaProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class.getName());
return new DefaultKafkaProducerFactory<>(kafkaProps);
}
}

View file

@ -104,6 +104,8 @@ public interface Consts {
String SDF04_PROCESS = "sdf04-process";
String SDF03_PROCESS = "sdf03-process";
String SDF11_PROCESS = "sdf11-process";
String SDF56_PROCESS = "sdf56-process";
String SDF57_PROCESS = "sdf57-process";
String EXPORT_PROCESS = "export-process";
String S_TRADES_IMPORTED = "s_trades-imported";
String LIM_EXPORTED = "lim_exported";

View file

@ -0,0 +1,38 @@
package ru.spcex.clearing.platform.messaging.domain.cud.clearing;
import com.fasterxml.jackson.annotation.JsonProperty;
import java.time.Instant;
public class Sdf56Request {
@JsonProperty
private Long number;
@JsonProperty
private Instant sDateTime;
@JsonProperty
private Instant eDateTime;
public Long getNumber() {
return number;
}
public void setNumber(Long number) {
this.number = number;
}
public Instant getsDateTime() {
return sDateTime;
}
public void setsDateTime(Instant sDateTime) {
this.sDateTime = sDateTime;
}
public Instant geteDateTime() {
return eDateTime;
}
public void seteDateTime(Instant eDateTime) {
this.eDateTime = eDateTime;
}
}

View file

@ -0,0 +1,18 @@
package ru.spcex.clearing.platform.messaging.serialization;
import ru.spcex.platform.utils.enumeration.IEnumKey;
public enum ClearingKafkaHeader implements IEnumKey {
BaseRequestTypeParameter("_cls_req_type_");
private final String key;
ClearingKafkaHeader(String key) {
this.key = key;
}
@Override
public String getKey() {
return key;
}
}

View file

@ -0,0 +1,59 @@
package ru.spcex.clearing.platform.messaging.serialization;
import com.fasterxml.jackson.databind.JavaType;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.kafka.common.header.Header;
import org.apache.kafka.common.header.Headers;
import org.apache.kafka.common.serialization.Deserializer;
import org.springframework.kafka.support.mapping.DefaultJackson2JavaTypeMapper;
import org.springframework.kafka.support.mapping.Jackson2JavaTypeMapper;
import org.springframework.util.ClassUtils;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
public class JsonDeserializer implements Deserializer<Object> {
private final ObjectMapper json;
private final Jackson2JavaTypeMapper typeMapper;
public JsonDeserializer() {
this.json = new ObjectMapper();
this.typeMapper = new DefaultJackson2JavaTypeMapper();
this.typeMapper.addTrustedPackages("*");
}
@Override
public Object deserialize(String topic, byte[] data) {
throw new RuntimeException("couldn't serialize kafka message without headers (need type info)");
}
@Override
public Object deserialize(String topic, Headers headers, byte[] data) {
try {
// ObjectReader deserReader = null;
// JavaType javaType = this.typeMapper.toJavaType(headers);
// if (javaType != null) {
// deserReader = this.json.readerFor(javaType);
// }
// if (deserReader == null) {
// throw new RuntimeException("couldn't serialize kafka message without headers (need type info)");
// }
Header header = headers.lastHeader(ClearingKafkaHeader.BaseRequestTypeParameter.getKey());
if (header == null) {
throw new RuntimeException("couldn't serialize kafka message without header '" +
ClearingKafkaHeader.BaseRequestTypeParameter.getKey() + "' (need type info)");
}
String typeParameterClassName = new String(header.value(), StandardCharsets.UTF_8);
// JavaType parameterType = json.getTypeFactory().constructType(ClassUtils.forName(typeParameterClassName, ClassUtils.getDefaultClassLoader()));
Class<?> parameterType = ClassUtils.forName(typeParameterClassName, ClassUtils.getDefaultClassLoader());
JavaType requestType = json.getTypeFactory().constructParametricType(BaseRequest.class, parameterType);
//...
headers.remove(ClearingKafkaHeader.BaseRequestTypeParameter.getKey());
return json.readValue(new String(data, StandardCharsets.UTF_8), requestType);
} catch (IOException | ClassNotFoundException e) {
throw new RuntimeException("couldn't deserialize kafka message ", e);
}
}
}

View file

@ -1,15 +1,20 @@
package ru.spcex.clearing.platform.messaging.serialization;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.kafka.common.header.Headers;
import org.apache.kafka.common.header.internals.RecordHeader;
import org.apache.kafka.common.serialization.Serializer;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import java.nio.charset.StandardCharsets;
public class JsonSerializer implements Serializer<Object> {
private final ObjectMapper json;
// private final Jackson2JavaTypeMapper typeMapper;
public JsonSerializer() {
this.json = new ObjectMapper();
// this.typeMapper = new DefaultJackson2JavaTypeMapper();
}
@Override
@ -21,4 +26,23 @@ public class JsonSerializer implements Serializer<Object> {
throw new RuntimeException("couldn't serialize kafka message ", e);
}
}
@Override
public byte[] serialize(String topic, Headers headers, Object data) {
if (data == null) {
return null;
}
//if (this.addTypeInfo && headers != null) {
//Class<?> clazz = callback.getClazz();
//JavaType payloadType = json.getTypeFactory().constructParametricType(BaseRequest.class, clazz);
//JavaType payloadType = json.getTypeFactory().constructParametricType(BaseRequest.class, clazz);
//JavaType payloadType = json.getTypeFactory().constructCollectionLikeType(BaseRequest.class, clazz);
//failed to use Jackson2JavaTypeMapper with parametrized BaseRequest
//this.typeMapper.fromJavaType(payloadType, headers);
BaseRequest<?> req = (BaseRequest<?>) data;
Class<?> clazz = req.getRequestPayload().getClass();
headers.add(new RecordHeader(ClearingKafkaHeader.BaseRequestTypeParameter.getKey(), clazz.getName().getBytes(StandardCharsets.UTF_8)));
return serialize(topic, data);
// }
}
}

View file

@ -1,10 +1,15 @@
package ru.spcex.clearing.platform.messaging.service.sender;
import org.apache.kafka.clients.consumer.ConsumerRecord;
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 org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.requestreply.ReplyingKafkaTemplate;
import org.springframework.kafka.requestreply.RequestReplyFuture;
import org.springframework.kafka.support.SendResult;
import org.springframework.util.concurrent.ListenableFuture;
import ru.spcex.clearing.platform.messaging.domain.ActionType;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
@ -14,13 +19,15 @@ 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.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
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 KafkaTemplate<String, Object> kafkaTemplate;
private Producer<String, Object> kafka;
private Supplier<Long> idGenerator;
private Function<String, KafkaImdgInsert> imdgGenerator;
@ -41,11 +48,11 @@ public class KafkaSender {
request.setRequestPayload(requestPayload);
//сохраняет данные о запросе в хранилище
saveRequestToStorage(destination, request);
ProducerRecord<Object, BaseRequest<Object>> prdRec = new ProducerRecord<>(destination, request);
ProducerRecord<String, Object> prdRec = new ProducerRecord<>(destination, request);
if (correlationId != null) {
prdRec.headers().add(new CorrelationHeader(correlationId.clone()));
}
Future<RecordMetadata> send = kafka.send(new ProducerRecord<>(destination, request));
ListenableFuture<SendResult<String, Object>> send = kafkaTemplate.send(prdRec);
try {
send.get();
} catch (InterruptedException | ExecutionException e) {
@ -58,6 +65,39 @@ public class KafkaSender {
return request.getId();
}
/**
* Вариант метода для случая, когда хотим додждаться ответа от другого сервиса через кафку
*/
public <R> R sendToQueueWaitForAnswer(String destination, Object requestPayload) {
if (!(kafkaTemplate instanceof ReplyingKafkaTemplate<?,?,?>)) {
throw new IllegalStateException("kafkaTemplate is not ReplyingKafkaTemplate");
}
BaseRequest<Object> request = new BaseRequest<>();
request.setId(idGenerator.get());
request.setActionType(ActionType.SYSTEM);
request.setRequestPayload(requestPayload);
//сохраняет данные о запросе в хранилище
saveRequestToStorage(destination, request);
ProducerRecord<String, Object> prdRec = new ProducerRecord<>(destination, request);
ReplyingKafkaTemplate<String, Object, Object> rplKafkaTemplate = (ReplyingKafkaTemplate<String, Object, Object>) kafkaTemplate;
RequestReplyFuture<String, Object, Object> replyFuture = rplKafkaTemplate.sendAndReceive(prdRec);
//can be used to extract Producer metadata
//SendResult<String, Object> sendResult = replyFuture.getSendFuture().get(10, TimeUnit.SECONDS);
Object value = null;
try {
ConsumerRecord<String, Object> consumerRecord = replyFuture.get(10, TimeUnit.SECONDS);
value = consumerRecord.value();
return (R) value;
} catch (InterruptedException | ExecutionException | TimeoutException | ClassCastException e) {
if (e instanceof InterruptedException) {
Thread.currentThread().interrupt();
}
log.error(ExceptionUtils.getStackTrace(e));
return null;
}
}
public Long sendRequestToQueue(String destination, Object requestPayload) {
return sendRequestToQueue(destination, requestPayload, null);
@ -80,8 +120,9 @@ public class KafkaSender {
return allImdgMaps.computeIfAbsent(mapName, (mapName1) -> imdgGenerator.apply(mapName));
}
@Deprecated
void setProducer(Producer<String, Object> kafka) {
this.kafka = kafka;
// this.kafka = kafka;
}
void setIdGenerator(Supplier<Long> idGenerator) {
@ -95,4 +136,8 @@ public class KafkaSender {
void setSaveRequestInfo(boolean saveRequestInfo) {
this.saveRequestInfo = saveRequestInfo;
}
void setKafkaTemplate(KafkaTemplate<String, Object> kafkaTemplate) {
this.kafkaTemplate = kafkaTemplate;
}
}

View file

@ -1,6 +1,7 @@
package ru.spcex.clearing.platform.messaging.service.sender;
import org.apache.kafka.clients.producer.Producer;
import org.springframework.kafka.core.KafkaTemplate;
import java.util.function.Function;
import java.util.function.Supplier;
@ -8,6 +9,8 @@ import java.util.function.Supplier;
public interface KafkaSenderBuilder {
KafkaSenderBuilder producer(Producer<String, Object> kafka);
KafkaSenderBuilder setKafkaTemplate(KafkaTemplate<String, Object> kafkaTemplate);
KafkaSenderBuilder idGenerator(Supplier<Long> idGenerator);
KafkaSenderBuilder imdgProvider(Function<String, KafkaImdgInsert> imdgGenerator);

View file

@ -1,6 +1,7 @@
package ru.spcex.clearing.platform.messaging.service.sender;
import org.apache.kafka.clients.producer.Producer;
import org.springframework.kafka.core.KafkaTemplate;
import java.util.function.Function;
import java.util.function.Supplier;
@ -19,6 +20,12 @@ public class KafkaSenderBuilderImpl implements KafkaSenderBuilder{
return this;
}
@Override
public KafkaSenderBuilder setKafkaTemplate(KafkaTemplate<String, Object> kafkaTemplate) {
kafkaSender.setKafkaTemplate(kafkaTemplate);
return this;
}
@Override
public KafkaSenderBuilder idGenerator(Supplier<Long> idGenerator) {
kafkaSender.setIdGenerator(idGenerator);