From 6c23cbca36833ab88d9a37863686e2505d881075 Mon Sep 17 00:00:00 2001 From: ialbert Date: Mon, 15 May 2023 12:04:05 +0300 Subject: [PATCH] =?UTF-8?q?http://jira.mfd.msk:8088/browse/CLS-290=20?= =?UTF-8?q?=D0=B4=D0=BE=D0=B1=D0=B0=D0=B2=D0=B8=D0=BB=20Spring=20KafkaTemp?= =?UTF-8?q?late?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../clearing/account/config/KafkaConfig.java | 18 +++++- .../clearing/balance/config/KafkaConfig.java | 18 +++++- .../ru/spcex/clearing/config/KafkaConfig.java | 22 +++++-- .../ru/spcex/clearing/service/Clearing.java | 3 +- .../clearing/company/config/KafkaConfig.java | 36 +++++++++++ .../company/config/KafkaSenderConfig.java | 32 ---------- .../exporter/config/KafkaSenderConfig.java | 36 ++++++++--- .../dbf/importer/config/KafkaConfig.java | 41 ++++++++----- .../iml/hazelcast/adapter/ImdgHazelcast.java | 31 ++++++++++ platform-parent/platform-imdg-api/pom.xml | 4 ++ .../java/ru/spcex/platform/imdg/api/Imdg.java | 28 +++++++++ .../specific/StatementRevisePredicate.java | 18 ++++++ .../config/KafkaConsumerFactory.java | 26 ++++++++ .../config/KafkaProducerFactory.java | 29 +++++++++ .../platform/messaging/domain/Consts.java | 2 + .../domain/cud/clearing/Sdf56Request.java | 38 ++++++++++++ .../serialization/ClearingKafkaHeader.java | 18 ++++++ .../serialization/JsonDeserializer.java | 59 +++++++++++++++++++ .../serialization/JsonSerializer.java | 24 ++++++++ .../messaging/service/sender/KafkaSender.java | 55 +++++++++++++++-- .../service/sender/KafkaSenderBuilder.java | 3 + .../sender/KafkaSenderBuilderImpl.java | 7 +++ 22 files changed, 478 insertions(+), 70 deletions(-) delete mode 100644 clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/KafkaSenderConfig.java create mode 100644 platform-parent/platform-imdg-api/src/main/java/ru/spcex/platform/imdg/api/predicate/specific/StatementRevisePredicate.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/Sdf56Request.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/serialization/ClearingKafkaHeader.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/serialization/JsonDeserializer.java diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/config/KafkaConfig.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/config/KafkaConfig.java index 58a707200..c57a09d19 100644 --- a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/config/KafkaConfig.java +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/config/KafkaConfig.java @@ -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 pf(AccountServiceSettings settings) { + KafkaProducerSettings kafkaSettings = settings.getKafkaProducer(); + return KafkaProducerFactory.producerFactory(kafkaSettings); + } + + @Bean + public KafkaTemplate kafkaTemplate(ProducerFactory pf) { + return new KafkaTemplate<>(pf); + } + @Autowired @Bean - public KafkaSender kafkaSender(Producer kafkaProducer, ImdgProvider imdgProvider) { + public KafkaSender kafkaSender(KafkaTemplate kafkaTemplate, ImdgProvider imdgProvider) { ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); return KafkaSender .setup() - .producer(kafkaProducer) + .setKafkaTemplate(kafkaTemplate) .idGenerator(imdgIdGenerator::nextId) .imdgProvider(s -> { Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/KafkaConfig.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/KafkaConfig.java index 1e14bcf54..ade27742a 100644 --- a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/KafkaConfig.java +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/KafkaConfig.java @@ -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 pf(BalanceServiceSettings settings) { + KafkaProducerSettings kafkaSettings = settings.getKafkaProducer(); + return KafkaProducerFactory.producerFactory(kafkaSettings); + } + + @Bean("kafkaTemplate") + public KafkaTemplate kafkaTemplate(ProducerFactory pf) { + return new KafkaTemplate<>(pf); + } + @Autowired @Bean - public KafkaSender kafkaSender(@Qualifier("kafkaProducer") Producer kafkaProducer, + public KafkaSender kafkaSender(@Qualifier("kafkaTemplate") KafkaTemplate kafkaTemplate, ImdgProvider imdgProvider) { ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); return KafkaSender .setup() - .producer(kafkaProducer) + .setKafkaTemplate(kafkaTemplate) .idGenerator(imdgIdGenerator::nextId) .imdgProvider(s -> { Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/KafkaConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/KafkaConfig.java index 5a7e0add1..95ad7f18b 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/KafkaConfig.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/KafkaConfig.java @@ -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 pf(ClearingServiceSettings settings) { + KafkaProducerSettings kafkaSettings = settings.getKafkaProducer(); + return KafkaProducerFactory.producerFactory(kafkaSettings); + } + + @Bean("kafkaTemplate") + public KafkaTemplate kafkaTemplate(ProducerFactory pf) { + return new KafkaTemplate<>(pf); + } + @Autowired @Bean - public KafkaSender kafkaSender(Producer kafkaProducer, ImdgProvider imdgProvider) { + public KafkaSender kafkaSender(KafkaTemplate kafkaTemplate, ImdgProvider imdgProvider) { ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); return KafkaSender .setup() - .producer(kafkaProducer) + .setKafkaTemplate(kafkaTemplate) .idGenerator(imdgIdGenerator::nextId) .imdgProvider(s -> { Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); @@ -49,11 +63,11 @@ public class KafkaConfig { @Autowired @Bean("kafkaSenderWithoutRequestInfo") - public KafkaSender kafkaSenderWithoutRequestInfo(Producer kafkaProducer, ImdgProvider imdgProvider) { + public KafkaSender kafkaSenderWithoutRequestInfo(KafkaTemplate kafkaTemplate, ImdgProvider imdgProvider) { ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); return KafkaSender .setup() - .producer(kafkaProducer) + .setKafkaTemplate(kafkaTemplate) .idGenerator(imdgIdGenerator::nextId) .saveRequestInfo(false) .build(); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Clearing.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Clearing.java index 1f8f795e0..0fabeca75 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Clearing.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Clearing.java @@ -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 validation, KafkaSender kafka, LiabilitiesClaimsAssetsCreator liabilitiesClaimsAssetsCreator, LiabilitiesClaimsMoneyCreator lbltsClmsMoneyCreator) { + BiFunction 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); diff --git a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/KafkaConfig.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/KafkaConfig.java index 40334c15c..1c99da87f 100644 --- a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/KafkaConfig.java +++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/KafkaConfig.java @@ -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 pf(CompanyServiceSettings settings) { + KafkaProducerSettings kafkaSettings = settings.getKafkaProducer(); + return KafkaProducerFactory.producerFactory(kafkaSettings); + } + + @Bean("kafkaTemplate") + public KafkaTemplate kafkaTemplate(ProducerFactory pf) { + return new KafkaTemplate<>(pf); + } + + @Autowired + @Bean + public KafkaSender 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/company-service/src/main/java/ru/spcex/clearing/company/config/KafkaSenderConfig.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/KafkaSenderConfig.java deleted file mode 100644 index 16e63b03d..000000000 --- a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/KafkaSenderConfig.java +++ /dev/null @@ -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 kafkaProducer, - ImdgProvider imdgProvider) { - ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); - return KafkaSender - .setup() - .producer(kafkaProducer) - .idGenerator(imdgIdGenerator::nextId) - .imdgProvider(s -> { - Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); - return imdg::insert; - }) - .build(); - } -} diff --git a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/KafkaSenderConfig.java b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/KafkaSenderConfig.java index 2b2319156..5f0a59e40 100644 --- a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/KafkaSenderConfig.java +++ b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/KafkaSenderConfig.java @@ -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 kafkaProducer; @Autowired public KafkaSenderConfig(ImdgProvider imdgProvider) { this.imdgProvider = imdgProvider; } - @Autowired(required = false) - public void setKafkaProducer(Producer kafkaProducer) { - this.kafkaProducer = kafkaProducer; + @Bean + public ProducerFactory pf(ExportDBFServiceSettings settings) { + if (settings.getKafkaProducer() == null) { + return null; + } + KafkaProducerSettings kafkaSettings = settings.getKafkaProducer(); + return KafkaProducerFactory.producerFactory(kafkaSettings); } + @Autowired(required = false) + @Bean("kafkaTemplate") + public KafkaTemplate kafkaTemplate(ProducerFactory 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 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 imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/KafkaConfig.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/KafkaConfig.java index fc448445b..398e6a9c5 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/KafkaConfig.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/KafkaConfig.java @@ -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 pf(ImportDBFServiceSettings settings) { + KafkaProducerSettings kafkaSettings = settings.getKafkaProducer(); + return KafkaProducerFactory.producerFactory(kafkaSettings); + } + + @Bean("kafkaTemplate") + public KafkaTemplate kafkaTemplate(ProducerFactory pf) { + return new KafkaTemplate<>(pf); + } + @Autowired @Bean - public Supplier kafkaSender(ImportDBFServiceSettings settings, ImdgProvider imdgProvider) { - return () -> { - Producer kafkaProducer = KafkaProducerFactory.producer(settings.getKafkaProducer()); - ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); - return KafkaSender - .setup() - .producer(kafkaProducer) - .idGenerator(imdgIdGenerator::nextId) - .imdgProvider(s -> { - Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); - return imdg::insert; - }) - .build(); - }; + 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(); } @Autowired diff --git a/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/adapter/ImdgHazelcast.java b/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/adapter/ImdgHazelcast.java index b087dae68..a544c3a45 100644 --- a/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/adapter/ImdgHazelcast.java +++ b/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/adapter/ImdgHazelcast.java @@ -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 implements Imdg { return map.values(); } + @Override + public Long aggregateLongMax(String fieldName, ImdgPredicate predicate) { + Predicate p = ((ImdgPredicateHazelcast) predicate).getRawPredicate(); + return map.aggregate(Aggregators.longMax(fieldName), p); + } + + @Override + public T aggregateByMax(String fieldName, ImdgPredicate predicate) { + Map.Entry entry; + if (predicate != null) { + Predicate 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 Collection projectSingleAttribute(String fieldName, ImdgPredicate predicate) { + if (predicate != null) { + Predicate 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(); diff --git a/platform-parent/platform-imdg-api/pom.xml b/platform-parent/platform-imdg-api/pom.xml index 678e60d10..1c6e12c22 100644 --- a/platform-parent/platform-imdg-api/pom.xml +++ b/platform-parent/platform-imdg-api/pom.xml @@ -25,6 +25,10 @@ ru.spcex.platform platform-classes-base + + ru.spcex.platform + platform-enum + ru.spcex.platform platform-utils diff --git a/platform-parent/platform-imdg-api/src/main/java/ru/spcex/platform/imdg/api/Imdg.java b/platform-parent/platform-imdg-api/src/main/java/ru/spcex/platform/imdg/api/Imdg.java index 8002b554b..68007e757 100644 --- a/platform-parent/platform-imdg-api/src/main/java/ru/spcex/platform/imdg/api/Imdg.java +++ b/platform-parent/platform-imdg-api/src/main/java/ru/spcex/platform/imdg/api/Imdg.java @@ -64,6 +64,34 @@ public interface Imdg { 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 Collection projectSingleAttribute(String fieldName, ImdgPredicate predicate) { + throw new UnsupportedOperationException("not implemented projectSingleAttribute"); + } + + /** + * при необходимости ввести ImdgProjection + */ + default Collection projectSingleAttribute(String fieldName) { + return projectSingleAttribute(fieldName, null); + } + default Integer size() { throw new UnsupportedOperationException("not implemented size"); } diff --git a/platform-parent/platform-imdg-api/src/main/java/ru/spcex/platform/imdg/api/predicate/specific/StatementRevisePredicate.java b/platform-parent/platform-imdg-api/src/main/java/ru/spcex/platform/imdg/api/predicate/specific/StatementRevisePredicate.java new file mode 100644 index 000000000..657a333f6 --- /dev/null +++ b/platform-parent/platform-imdg-api/src/main/java/ru/spcex/platform/imdg/api/predicate/specific/StatementRevisePredicate.java @@ -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 currenciesId) { + ImdgPredicateBuilder prdctBuilder = pb.predicateBuilder(); + return prdctBuilder.and( + prdctBuilder.in("currencyId", currenciesId.toArray(new Comparable[0])), + prdctBuilder.equals("statementType", StatementType.incr.getKey()) + ); + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaConsumerFactory.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaConsumerFactory.java index f2c8a7fa3..26e7209f8 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaConsumerFactory.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaConsumerFactory.java @@ -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 consumerFactory(KafkaConsumerSettings kafkaSettings) { + Map 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 cf = new DefaultKafkaConsumerFactory<>(props); + cf.setValueDeserializer(new ErrorHandlingDeserializer<>(new JsonDeserializer())); + return cf; + } } diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaProducerFactory.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaProducerFactory.java index 059dfb381..fdae5c834 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaProducerFactory.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaProducerFactory.java @@ -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 producerFactory(KafkaProducerSettings kafkaSettings) { + Map 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); + } } diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java index 6523560bc..1991be59e 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java @@ -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"; diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/Sdf56Request.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/Sdf56Request.java new file mode 100644 index 000000000..2dd1eb75e --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/Sdf56Request.java @@ -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; + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/serialization/ClearingKafkaHeader.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/serialization/ClearingKafkaHeader.java new file mode 100644 index 000000000..69cd78760 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/serialization/ClearingKafkaHeader.java @@ -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; + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/serialization/JsonDeserializer.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/serialization/JsonDeserializer.java new file mode 100644 index 000000000..ca0ceb998 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/serialization/JsonDeserializer.java @@ -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 { + + 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); + } + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/serialization/JsonSerializer.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/serialization/JsonSerializer.java index 5e9cedf82..e819ab119 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/serialization/JsonSerializer.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/serialization/JsonSerializer.java @@ -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 { 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 { 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); +// } + } } diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSender.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSender.java index 35617813f..39ec6b504 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSender.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSender.java @@ -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 allImdgMaps; + private KafkaTemplate kafkaTemplate; private Producer kafka; private Supplier idGenerator; private Function imdgGenerator; @@ -41,11 +48,11 @@ public class KafkaSender { request.setRequestPayload(requestPayload); //сохраняет данные о запросе в хранилище saveRequestToStorage(destination, request); - ProducerRecord> prdRec = new ProducerRecord<>(destination, request); + ProducerRecord prdRec = new ProducerRecord<>(destination, request); if (correlationId != null) { prdRec.headers().add(new CorrelationHeader(correlationId.clone())); } - Future send = kafka.send(new ProducerRecord<>(destination, request)); + ListenableFuture> send = kafkaTemplate.send(prdRec); try { send.get(); } catch (InterruptedException | ExecutionException e) { @@ -58,6 +65,39 @@ public class KafkaSender { return request.getId(); } + /** + * Вариант метода для случая, когда хотим додждаться ответа от другого сервиса через кафку + */ + public R sendToQueueWaitForAnswer(String destination, Object requestPayload) { + if (!(kafkaTemplate instanceof ReplyingKafkaTemplate)) { + throw new IllegalStateException("kafkaTemplate is not ReplyingKafkaTemplate"); + } + BaseRequest request = new BaseRequest<>(); + request.setId(idGenerator.get()); + request.setActionType(ActionType.SYSTEM); + request.setRequestPayload(requestPayload); + //сохраняет данные о запросе в хранилище + saveRequestToStorage(destination, request); + ProducerRecord prdRec = new ProducerRecord<>(destination, request); + ReplyingKafkaTemplate rplKafkaTemplate = (ReplyingKafkaTemplate) kafkaTemplate; + + RequestReplyFuture replyFuture = rplKafkaTemplate.sendAndReceive(prdRec); + //can be used to extract Producer metadata + //SendResult sendResult = replyFuture.getSendFuture().get(10, TimeUnit.SECONDS); + Object value = null; + try { + ConsumerRecord 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 kafka) { - this.kafka = kafka; +// this.kafka = kafka; } void setIdGenerator(Supplier idGenerator) { @@ -95,4 +136,8 @@ public class KafkaSender { void setSaveRequestInfo(boolean saveRequestInfo) { this.saveRequestInfo = saveRequestInfo; } + + void setKafkaTemplate(KafkaTemplate kafkaTemplate) { + this.kafkaTemplate = kafkaTemplate; + } } diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSenderBuilder.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSenderBuilder.java index 0a6d2f6af..2ebe9cc6d 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSenderBuilder.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSenderBuilder.java @@ -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 kafka); + KafkaSenderBuilder setKafkaTemplate(KafkaTemplate kafkaTemplate); + KafkaSenderBuilder idGenerator(Supplier idGenerator); KafkaSenderBuilder imdgProvider(Function imdgGenerator); diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSenderBuilderImpl.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSenderBuilderImpl.java index 0e89ad5cb..35306e0c5 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSenderBuilderImpl.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSenderBuilderImpl.java @@ -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 kafkaTemplate) { + kafkaSender.setKafkaTemplate(kafkaTemplate); + return this; + } + @Override public KafkaSenderBuilder idGenerator(Supplier idGenerator) { kafkaSender.setIdGenerator(idGenerator);