diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/BankAccountService.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/BankAccountService.java index b95313fd0..1293f4d9d 100644 --- a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/BankAccountService.java +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/BankAccountService.java @@ -1,6 +1,7 @@ package ru.spcex.clearing.account.service; import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; @@ -23,8 +24,8 @@ public class BankAccountService extends QueueConsumer implements InitializingBea private final Imdg bankAccountMap; @Autowired - public BankAccountService(Consumer kafkaQueue, ImdgProvider imdgProvider) { - super(kafkaQueue); + public BankAccountService(Consumer kafkaQueue, Producer kafkaProducer, ImdgProvider imdgProvider) { + super(kafkaQueue, kafkaProducer); this.bankAccountMap = imdgProvider.getImdg(IMDGDistributedNames.Map_BankAccount, BankAccount.class); } diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/KafkaConfig.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/KafkaConfig.java index f64045cb9..9046c9bc8 100644 --- a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/KafkaConfig.java +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/KafkaConfig.java @@ -1,5 +1,8 @@ package ru.spcex.clearing.backendapi.config; +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.consumer.MockConsumer; +import org.apache.kafka.clients.consumer.OffsetResetStrategy; import org.apache.kafka.clients.producer.MockProducer; import org.apache.kafka.clients.producer.Producer; import org.springframework.beans.factory.annotation.Autowired; @@ -7,6 +10,7 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Profile; import ru.spcex.clearing.backendapi.config.element.BackendApiSettings; +import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; @Configuration @@ -16,7 +20,7 @@ public class KafkaConfig { @Autowired @Bean public Producer createProducer(BackendApiSettings settings) { - return KafkaProducerFactory.producer(settings.getKafka()); + return KafkaProducerFactory.producer(settings.getKafkaProducer()); } @Profile("kafkaDisabled") @@ -25,4 +29,17 @@ public class KafkaConfig { return new MockProducer<>(true, (topic, data) -> data != null ? data.getBytes() : new byte[0], (topic, data) -> data.toString().getBytes()); } + @Profile("!kafkaDisabled") + @Autowired + @Bean + public Consumer createConsumer(BackendApiSettings settings) { + return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); + } + + @Profile("kafkaDisabled") + @Bean + public Consumer createConsumerMock() { + return new MockConsumer<>(OffsetResetStrategy.NONE); + } + } diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/element/BackendApiSettings.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/element/BackendApiSettings.java index 422d18323..79ba99f2d 100644 --- a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/element/BackendApiSettings.java +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/element/BackendApiSettings.java @@ -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.backendapi.security.element.SecuritySettings; +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; @@ -12,7 +13,8 @@ import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams; @ConfigurationProperties("backend-api") public class BackendApiSettings { private HazelcastClientParams hazelcast; - private KafkaProducerSettings kafka; + private KafkaProducerSettings kafkaProducer; + private KafkaConsumerSettings kafkaConsumer; private SecuritySettings security; private String exampleSetting; @@ -24,12 +26,20 @@ public class BackendApiSettings { this.hazelcast = hazelcast; } - public KafkaProducerSettings getKafka() { - return kafka; + public KafkaProducerSettings getKafkaProducer() { + return kafkaProducer; } - public void setKafka(KafkaProducerSettings kafka) { - this.kafka = kafka; + public void setKafkaProducer(KafkaProducerSettings kafkaProducer) { + this.kafkaProducer = kafkaProducer; + } + + public KafkaConsumerSettings getKafkaConsumer() { + return kafkaConsumer; + } + + public void setKafkaConsumer(KafkaConsumerSettings kafkaConsumer) { + this.kafkaConsumer = kafkaConsumer; } public String getExampleSetting() { diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/service/RequestInfoAccepter.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/service/RequestInfoAccepter.java new file mode 100644 index 000000000..a1eaff115 --- /dev/null +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/service/RequestInfoAccepter.java @@ -0,0 +1,48 @@ +package ru.spcex.clearing.backendapi.service; + +import org.apache.kafka.clients.consumer.Consumer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.stereotype.Service; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.clearing.platform.messaging.service.RequestInfo; +import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +@Service +public class RequestInfoAccepter extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg requestInfoImdg; + + + public RequestInfoAccepter(Consumer kafkaQueue, ImdgProvider imdgProvider) { + super(kafkaQueue); + this.requestInfoImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); + } + + + @Override + public void afterPropertiesSet() throws Exception { + callback(RequestInfoUpdate.class) + .setConsumer(this::updateRequestInfo) + .forDestination(Consts.REQUEST_INFO_UPDATE, callbacks::put); + init(); + } + + private void updateRequestInfo(BaseRequest requestInfoUpdateBaseRequest) { + //will throw exception for any class other than RequestInfoUpdate + RequestInfoUpdate statusInfo = requestInfoUpdateBaseRequest.getRequestPayload(); + RequestInfo requestInfo = requestInfoImdg.getSingleObjectByID(statusInfo.getId()); + if (requestInfo == null) { + log.warn("unknown requestInfo id={}", statusInfo.getId()); + return; + } + requestInfo.setStatus(statusInfo.getStatus()); + requestInfoImdg.update(requestInfo); + } +} diff --git a/clearing-parent/backend-api/src/main/resources/application.properties b/clearing-parent/backend-api/src/main/resources/application.properties index 32e6e67c6..c375b65ed 100644 --- a/clearing-parent/backend-api/src/main/resources/application.properties +++ b/clearing-parent/backend-api/src/main/resources/application.properties @@ -12,14 +12,22 @@ backend-api.hazelcast.cluster-members=127.0.0.1:5701 backend-api.hazelcast.login=dev backend-api.hazelcast.password=dev-pass -backend-api.kafka.bootstrap-servers=localhost:9092 -backend-api.kafka.acks=all -backend-api.kafka.retries=0 -backend-api.kafka.batch-size=16384 -backend-api.kafka.linger-ms=1 -backend-api.kafka.buffer-memory=33554432 +backend-api.kafka-producer.bootstrap-servers=localhost:9092 +backend-api.kafka-producer.acks=all +backend-api.kafka-producer.retries=0 +backend-api.kafka-producer.batch-size=16384 +backend-api.kafka-producer.linger-ms=1 +backend-api.kafka-producer.buffer-memory=33554432 -backend-api.security.authorization-disabled=true +backend-api.kafka-consumer.bootstrap-servers=localhost:9092 +backend-api.kafka-consumer.group-id=dev-group-backend-api +backend-api.kafka-consumer.enable-auto-commit=true +backend-api.kafka-consumer.session-timeout-ms=30000 +backend-api.kafka-consumer.auto-offset-reset=latest +backend-api.kafka-consumer.linger-ms=1 +backend-api.kafka-consumer.buffer-memory=33554432 + +backend-api.security.authorization-disabled=false ##keycloak 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 6ce4bca40..0be729887 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 @@ -1,6 +1,7 @@ package ru.spcex.clearing.company.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,6 +9,7 @@ import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Scope; import ru.spcex.clearing.company.config.settings.CompanyServiceSettings; import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; +import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; @Configuration public class KafkaConfig { @@ -15,6 +17,13 @@ public class KafkaConfig { @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) @Bean public Consumer createConsumer(CompanyServiceSettings settings) { - return KafkaConsumerFactory.consumer(settings.getKafka()); + return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); } + + @Autowired + @Bean + public Producer createProducer(CompanyServiceSettings settings) { + return KafkaProducerFactory.producer(settings.getKafkaProducer()); + } + } diff --git a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/settings/CompanyServiceSettings.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/settings/CompanyServiceSettings.java index e67bb2f8c..c27eae7f5 100644 --- a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/settings/CompanyServiceSettings.java +++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/settings/CompanyServiceSettings.java @@ -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("company-service") public class CompanyServiceSettings { private HazelcastClientParams hazelcast; - private KafkaConsumerSettings kafka; + private KafkaConsumerSettings kafkaConsumer; + private KafkaProducerSettings kafkaProducer; public HazelcastClientParams getHazelcast() { return hazelcast; @@ -21,11 +23,19 @@ public class CompanyServiceSettings { 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; } } diff --git a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/ClearingMemberCategoryService.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/ClearingMemberCategoryService.java index 7ed87f27b..2cfab579f 100644 --- a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/ClearingMemberCategoryService.java +++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/ClearingMemberCategoryService.java @@ -1,6 +1,7 @@ package ru.spcex.clearing.company.service; import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; @@ -20,9 +21,9 @@ public class ClearingMemberCategoryService extends QueueConsumer implements Init private final Imdg clearingMemberCategoryMap; @Autowired - public ClearingMemberCategoryService(Consumer kafkaQueue, + public ClearingMemberCategoryService(Consumer kafkaQueue, Producer kafkaProducer, ImdgProvider imdgProvider) { - super(kafkaQueue); + super(kafkaQueue, kafkaProducer); this.clearingMemberCategoryMap = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingMemberCategory, ClearingMemberCategory.class); } diff --git a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/CompanyInfoService.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/CompanyInfoService.java index 85775c9bb..953bd35b8 100644 --- a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/CompanyInfoService.java +++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/CompanyInfoService.java @@ -1,6 +1,7 @@ package ru.spcex.clearing.company.service; import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; @@ -23,9 +24,9 @@ public class CompanyInfoService extends QueueConsumer implements InitializingBea private final Imdg companyMap; @Autowired - public CompanyInfoService(Consumer kafkaQueue, + public CompanyInfoService(Consumer kafkaQueue, Producer kafkaProducer, ImdgProvider imdgProvider) { - super(kafkaQueue); + super(kafkaQueue, kafkaProducer); // this.companyInfoMap = imdgProvider.getImdg(IMDGDistributedNames.Map_CompanyInfo, CompanyInfo.class); this.companyMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class); } diff --git a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/CompanyService.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/CompanyService.java index 5a96460ac..e9008c33d 100644 --- a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/CompanyService.java +++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/CompanyService.java @@ -1,6 +1,7 @@ package ru.spcex.clearing.company.service; import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; @@ -21,9 +22,9 @@ public class CompanyService extends QueueConsumer implements InitializingBean { private final Imdg companyMap; @Autowired - public CompanyService(Consumer kafkaQueue, + public CompanyService(Consumer kafkaQueue, Producer kafkaProducer, ImdgProvider imdgProvider) { - super(kafkaQueue); + super(kafkaQueue, kafkaProducer); this.companyMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class); } diff --git a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/CompanySymbolService.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/CompanySymbolService.java index 4a6f110af..152fac925 100644 --- a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/CompanySymbolService.java +++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/CompanySymbolService.java @@ -2,6 +2,7 @@ package ru.spcex.clearing.company.service; import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; @@ -22,9 +23,9 @@ public class CompanySymbolService extends QueueConsumer implements InitializingB private final Imdg companySymbolsMap; @Autowired - public CompanySymbolService(Consumer kafkaQueue, + public CompanySymbolService(Consumer kafkaQueue, Producer kafkaProducer, ImdgProvider imdgProvider) { - super(kafkaQueue); + super(kafkaQueue, kafkaProducer); this.companySymbolsMap = imdgProvider.getImdg(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class); } diff --git a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/ContactService.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/ContactService.java index 0decc85f3..b38bf9ab2 100644 --- a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/ContactService.java +++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/ContactService.java @@ -1,6 +1,7 @@ package ru.spcex.clearing.company.service; import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; @@ -19,9 +20,9 @@ public class ContactService extends QueueConsumer implements InitializingBean { private final Imdg contactMap; @Autowired - public ContactService(Consumer kafkaQueue, + public ContactService(Consumer kafkaQueue, Producer kafkaProducer, ImdgProvider imdgProvider) { - super(kafkaQueue); + super(kafkaQueue, kafkaProducer); this.contactMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Contact, Contact.class); } diff --git a/clearing-parent/company-service/src/main/resources/application.properties b/clearing-parent/company-service/src/main/resources/application.properties index 292578a7e..6317c6c42 100644 --- a/clearing-parent/company-service/src/main/resources/application.properties +++ b/clearing-parent/company-service/src/main/resources/application.properties @@ -2,10 +2,17 @@ spring.main.web-application-type=none company-service.hazelcast.cluster-members=127.0.0.1:5701 company-service.hazelcast.login=dev company-service.hazelcast.password=dev-pass -company-service.kafka.bootstrap-servers=localhost:9092 -company-service.kafka.group-id=dev-group-company-service -company-service.kafka.enable-auto-commit=true -company-service.kafka.session-timeout-ms=30000 -company-service.kafka.auto-offset-reset=latest -company-service.kafka.linger-ms=1 -company-service.kafka.buffer-memory=33554432 \ No newline at end of file +company-service.kafka-consumer.bootstrap-servers=localhost:9092 +company-service.kafka-consumer.group-id=dev-group-company-service +company-service.kafka-consumer.enable-auto-commit=true +company-service.kafka-consumer.session-timeout-ms=30000 +company-service.kafka-consumer.auto-offset-reset=latest +company-service.kafka-consumer.linger-ms=1 +company-service.kafka-consumer.buffer-memory=33554432 + +company-service.kafka-producer.bootstrap-servers=localhost:9092 +company-service.kafka-producer.acks=all +company-service.kafka-producer.retries=0 +company-service.kafka-producer.batch-size=16384 +company-service.kafka-producer.linger-ms=1 +company-service.kafka-producer.buffer-memory=33554432 \ No newline at end of file diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/KafkaConfig.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/KafkaConfig.java index deb72e8ca..1669d0e82 100644 --- a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/KafkaConfig.java +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/KafkaConfig.java @@ -1,10 +1,12 @@ package ru.spcex.clearing.securities.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.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; +import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; import ru.spcex.clearing.securities.config.element.SecuritiesServiceSettings; @Configuration @@ -12,8 +14,15 @@ public class KafkaConfig { @Autowired @Bean - public Consumer createProducer(SecuritiesServiceSettings settings) { - return KafkaConsumerFactory.consumer(settings.getKafka()); + public Consumer createConsumer(SecuritiesServiceSettings settings) { + return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); } + @Autowired + @Bean + public Producer createProducer(SecuritiesServiceSettings settings) { + return KafkaProducerFactory.producer(settings.getKafkaProducer()); + } + + } diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/element/SecuritiesServiceSettings.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/element/SecuritiesServiceSettings.java index 546f25cb6..2ef7db60d 100644 --- a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/element/SecuritiesServiceSettings.java +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/element/SecuritiesServiceSettings.java @@ -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("securities-service") public class SecuritiesServiceSettings { private HazelcastClientParams hazelcast; - private KafkaConsumerSettings kafka; + private KafkaConsumerSettings kafkaConsumer; + private KafkaProducerSettings kafkaProducer; public HazelcastClientParams getHazelcast() { return hazelcast; @@ -21,11 +23,19 @@ public class SecuritiesServiceSettings { 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; } } diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java index 279c09bd5..a195fc683 100644 --- a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java @@ -1,6 +1,7 @@ package ru.spcex.clearing.securities.service.cud; import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; @@ -19,7 +20,6 @@ import ru.spcex.clearing.platform.messaging.service.Status; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; import ru.spcex.platform.imdg.api.ImdgTransaction; -import ru.spcex.platform.utils.log.ExceptionUtils; import java.math.BigDecimal; @@ -29,15 +29,16 @@ public class MoneyMarketSecurityService extends QueueConsumer implements Initial private final Imdg moneyMarketSecurityMap; private final ImdgProvider imdgProvider; @Autowired - public MoneyMarketSecurityService(Consumer kafkaQueue, ImdgProvider imdgProvider) { - super(kafkaQueue); + public MoneyMarketSecurityService(Consumer kafkaQueue, Producer kafkaProducer, + ImdgProvider imdgProvider) { + super(kafkaQueue, kafkaProducer); this.moneyMarketSecurityMap = imdgProvider.getImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, MoneyMarketSecurity.class); this.imdgProvider = imdgProvider; } @Override protected boolean needsProcessing(String destination, BaseRequest request) { - Imdg requestSpecificImdg = imdgProvider.getImdg(destination, RequestInfo.class); + Imdg requestSpecificImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); RequestInfo requestFound = requestSpecificImdg.getSingleObjectByID(request.getId()); if (requestFound == null) { log.error("couldn't find request info by id={}, action type={}", request.getId(), request.getActionType()); @@ -66,8 +67,8 @@ public class MoneyMarketSecurityService extends QueueConsumer implements Initial } private void newMoneyMarket(BaseRequest userRequest) { - Imdg requestInfoImdg = imdgProvider.getImdg(Consts.DESTINATION_MONEY_MARKET_SECURITY_NEW, RequestInfo.class); - RequestInfo reqInfo = requestInfoImdg.getSingleObjectByID(userRequest.getId()); +// Imdg requestInfoImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); +// RequestInfo reqInfo = requestInfoImdg.getSingleObjectByID(userRequest.getId()); ImdgTransaction transaction = imdgProvider.newTransaction(); MoneyMarketSecurityNewRequest req = userRequest.getRequestPayload(); log.debug("MoneyMarketSecurityNewRequest received"); @@ -82,17 +83,17 @@ public class MoneyMarketSecurityService extends QueueConsumer implements Initial try { transaction.beginTransaction(); Imdg moneyMarketSecurityMap = transaction.getImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, MoneyMarketSecurity.class); - Imdg reqInfoMap = transaction.getImdg(Consts.DESTINATION_MONEY_MARKET_SECURITY_NEW, RequestInfo.class); + Imdg reqInfoMap = transaction.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); //одним скопом выполняем реквест moneyMarketSecurityMap.insert(mms); - RequestInfo.update(reqInfo, Status.Success, "success"); +// RequestInfo.update(reqInfo, Status.Success, "success"); //и сохраняем обновленный - reqInfoMap.insert(reqInfo); +// reqInfoMap.insert(reqInfo); transaction.commitTransaction(); } catch (Throwable e) { transaction.rollbackTransaction(); - RequestInfo.update(reqInfo, Status.Error, ExceptionUtils.getStackTrace(e)); - requestInfoImdg.insert(reqInfo); +// RequestInfo.update(reqInfo, Status.Error, ExceptionUtils.getStackTrace(e)); +// requestInfoImdg.insert(reqInfo); } log.debug("successfully processed, request id {}, new object id {}", userRequest.getId(), mms.getId()); } diff --git a/clearing-parent/securities-service/src/main/resources/application.properties b/clearing-parent/securities-service/src/main/resources/application.properties index ee89c6710..45aa08e51 100644 --- a/clearing-parent/securities-service/src/main/resources/application.properties +++ b/clearing-parent/securities-service/src/main/resources/application.properties @@ -11,6 +11,13 @@ securities-service.kafka.session-timeout-ms=30000 securities-service.kafka.auto-offset-reset=latest securities-service.kafka.linger-ms=1 securities-service.kafka.buffer-memory=33554432 + +securities-service.kafka-producer.bootstrap-servers=localhost:9092 +securities-service.kafka-producer.acks=all +securities-service.kafka-producer.retries=0 +securities-service.kafka-producer.batch-size=16384 +securities-service.kafka-producer.linger-ms=1 +securities-service.kafka-producer.buffer-memory=33554432 #props.put("bootstrap.servers", kafkaSettings.getBootstrapServers()); #if (kafkaSettings.getGroupId() != null && kafkaSettings.getGroupId().length() > 0) { # props.put("group.id", kafkaSettings.getGroupId()); diff --git a/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/config/KafkaConfig.java b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/config/KafkaConfig.java index bc8a2546a..2dde7ea72 100644 --- a/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/config/KafkaConfig.java +++ b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/config/KafkaConfig.java @@ -1,12 +1,14 @@ package ru.spcex.clearing.utility.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 ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; +import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; import ru.spcex.clearing.utility.config.settings.UtilityServiceSettings; @Configuration @@ -14,7 +16,14 @@ public class KafkaConfig { @Autowired @Bean @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) - public Consumer createProducer(UtilityServiceSettings settings) { - return KafkaConsumerFactory.consumer(settings.getKafka()); + public Consumer createConsumer(UtilityServiceSettings settings) { + return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); } + + @Autowired + @Bean + public Producer createProducer(UtilityServiceSettings settings) { + return KafkaProducerFactory.producer(settings.getKafkaProducer()); + } + } diff --git a/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/config/settings/UtilityServiceSettings.java b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/config/settings/UtilityServiceSettings.java index 54f22e60c..f7db55b39 100644 --- a/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/config/settings/UtilityServiceSettings.java +++ b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/config/settings/UtilityServiceSettings.java @@ -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("utility-service") public class UtilityServiceSettings { private HazelcastClientParams hazelcast; - private KafkaConsumerSettings kafka; + private KafkaConsumerSettings kafkaConsumer; + private KafkaProducerSettings kafkaProducer; public HazelcastClientParams getHazelcast() { return hazelcast; @@ -21,11 +23,19 @@ public class UtilityServiceSettings { 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; } } diff --git a/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/KeyRateService.java b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/KeyRateService.java index 92c3cdf26..b96fa1032 100644 --- a/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/KeyRateService.java +++ b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/KeyRateService.java @@ -1,6 +1,7 @@ package ru.spcex.clearing.utility.service; import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; @@ -23,8 +24,9 @@ public class KeyRateService extends QueueConsumer implements InitializingBean { private final Imdg keyRateMap; @Autowired - public KeyRateService(Consumer kafkaQueue, ImdgProvider imdgProvider) { - super(kafkaQueue); + public KeyRateService(Consumer kafkaQueue, Producer kafkaProducer, + ImdgProvider imdgProvider) { + super(kafkaQueue, kafkaProducer); this.keyRateMap = imdgProvider.getImdg(IMDGDistributedNames.Map_KeyRate, KeyRate.class); } diff --git a/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/UserSettingService.java b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/UserSettingService.java index f34af7255..0c069f2fb 100644 --- a/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/UserSettingService.java +++ b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/UserSettingService.java @@ -1,6 +1,7 @@ package ru.spcex.clearing.utility.service; import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; @@ -23,8 +24,9 @@ public class UserSettingService extends QueueConsumer implements InitializingBea private final Imdg userSettingsMap; @Autowired - public UserSettingService(Consumer kafkaQueue, ImdgProvider imdgProvider) { - super(kafkaQueue); + public UserSettingService(Consumer kafkaQueue, Producer kafkaProducer, + ImdgProvider imdgProvider) { + super(kafkaQueue, kafkaProducer); this.userSettingsMap = imdgProvider.getImdg(IMDGDistributedNames.Map_UserSettings, UserSettings.class); } diff --git a/clearing-parent/utility-service/src/main/resources/application.properties b/clearing-parent/utility-service/src/main/resources/application.properties index fbc306eab..73a36eb86 100644 --- a/clearing-parent/utility-service/src/main/resources/application.properties +++ b/clearing-parent/utility-service/src/main/resources/application.properties @@ -4,10 +4,17 @@ utility-service.hazelcast.cluster-members=127.0.0.1:5701 utility-service.hazelcast.login=dev utility-service.hazelcast.password=dev-pass -utility-service.kafka.bootstrap-servers=localhost:9092 -utility-service.kafka.group-id=dev-group-utility-service -utility-service.kafka.enable-auto-commit=true -utility-service.kafka.session-timeout-ms=30000 -utility-service.kafka.auto-offset-reset=latest -utility-service.kafka.linger-ms=1 -utility-service.kafka.buffer-memory=33554432 \ No newline at end of file +utility-service.kafka-consumer.bootstrap-servers=localhost:9092 +utility-service.kafka-consumer.group-id=dev-group-utility-service +utility-service.kafka-consumer.enable-auto-commit=true +utility-service.kafka-consumer.session-timeout-ms=30000 +utility-service.kafka-consumer.auto-offset-reset=latest +utility-service.kafka-consumer.linger-ms=1 +utility-service.kafka-consumer.buffer-memory=33554432 + +utility-service.kafka-producer.bootstrap-servers=localhost:9092 +utility-service.kafka-producer.acks=all +utility-service.kafka-producer.retries=0 +utility-service.kafka-producer.batch-size=16384 +utility-service.kafka-producer.linger-ms=1 +utility-service.kafka-producer.buffer-memory=33554432 \ No newline at end of file diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java index da27d7857..0e11d3bb8 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java @@ -11,6 +11,7 @@ import org.apache.kafka.clients.producer.RecordMetadata; import org.apache.kafka.common.errors.WakeupException; 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.domain.Consts; import ru.spcex.clearing.platform.messaging.logic.functional.BuilderConsumerStep; @@ -105,10 +106,14 @@ public class QueueConsumer implements AutoCloseable { try { Future send; //default response + BaseRequest req = new BaseRequest<>(); RequestInfoUpdate success = new RequestInfoUpdate(); success.setId(o.getId()); success.setStatus(Status.Error); - send = producer.send(new ProducerRecord<>(Consts.REQUEST_INFO_UPDATE, success)); + req.setRequestPayload(success); + req.setId(o.getId()); + req.setActionType(ActionType.SYSTEM); + send = producer.send(new ProducerRecord<>(Consts.REQUEST_INFO_UPDATE, req)); send.get(); } catch (Exception e) { log.error(ExceptionUtils.getStackTrace(e)); @@ -120,15 +125,19 @@ public class QueueConsumer implements AutoCloseable { outputExecutor.submit(() -> { try { Future send; + BaseRequest req = new BaseRequest<>(); + req.setId(o.getId()); + req.setActionType(ActionType.SYSTEM); if (response != null) { - send = producer.send(new ProducerRecord<>(Consts.REQUEST_INFO_UPDATE, response)); + req.setRequestPayload(response); } else { //default response RequestInfoUpdate success = new RequestInfoUpdate(); success.setId(o.getId()); success.setStatus(Status.Success); - send = producer.send(new ProducerRecord<>(Consts.REQUEST_INFO_UPDATE, success)); + req.setRequestPayload(success); } + send = producer.send(new ProducerRecord<>(Consts.REQUEST_INFO_UPDATE, req)); send.get(); } catch (Exception e) { log.error(ExceptionUtils.getStackTrace(e));