ialbert 2022-10-12 13:17:21 +03:00
parent 5f2abbd2d9
commit ab407edb8c
23 changed files with 258 additions and 77 deletions

View file

@ -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<BankAccount> bankAccountMap;
@Autowired
public BankAccountService(Consumer<String, Object> kafkaQueue, ImdgProvider imdgProvider) {
super(kafkaQueue);
public BankAccountService(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaProducer, ImdgProvider imdgProvider) {
super(kafkaQueue, kafkaProducer);
this.bankAccountMap = imdgProvider.getImdg(IMDGDistributedNames.Map_BankAccount, BankAccount.class);
}

View file

@ -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<String, Object> 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<String, Object> createConsumer(BackendApiSettings settings) {
return KafkaConsumerFactory.consumer(settings.getKafkaConsumer());
}
@Profile("kafkaDisabled")
@Bean
public Consumer<String, Object> createConsumerMock() {
return new MockConsumer<>(OffsetResetStrategy.NONE);
}
}

View file

@ -4,6 +4,7 @@ import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.PropertySource;
import org.springframework.stereotype.Component;
import ru.spcex.clearing.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() {

View file

@ -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<RequestInfo> requestInfoImdg;
public RequestInfoAccepter(Consumer<String, Object> 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<RequestInfoUpdate> 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);
}
}

View file

@ -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

View file

@ -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<String, Object> createConsumer(CompanyServiceSettings settings) {
return KafkaConsumerFactory.consumer(settings.getKafka());
return KafkaConsumerFactory.consumer(settings.getKafkaConsumer());
}
@Autowired
@Bean
public Producer<String, Object> createProducer(CompanyServiceSettings settings) {
return KafkaProducerFactory.producer(settings.getKafkaProducer());
}
}

View file

@ -4,6 +4,7 @@ import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.PropertySource;
import org.springframework.stereotype.Component;
import ru.spcex.clearing.platform.messaging.config.element.KafkaConsumerSettings;
import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings;
import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams;
@Component
@ -11,7 +12,8 @@ import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams;
@ConfigurationProperties("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;
}
}

View file

@ -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<ClearingMemberCategory> clearingMemberCategoryMap;
@Autowired
public ClearingMemberCategoryService(Consumer<String, Object> kafkaQueue,
public ClearingMemberCategoryService(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider) {
super(kafkaQueue);
super(kafkaQueue, kafkaProducer);
this.clearingMemberCategoryMap = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingMemberCategory, ClearingMemberCategory.class);
}

View file

@ -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<Company> companyMap;
@Autowired
public CompanyInfoService(Consumer<String, Object> kafkaQueue,
public CompanyInfoService(Consumer<String, Object> kafkaQueue, Producer<String, Object> 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);
}

View file

@ -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<Company> companyMap;
@Autowired
public CompanyService(Consumer<String, Object> kafkaQueue,
public CompanyService(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider) {
super(kafkaQueue);
super(kafkaQueue, kafkaProducer);
this.companyMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
}

View file

@ -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<CompanySymbols> companySymbolsMap;
@Autowired
public CompanySymbolService(Consumer<String, Object> kafkaQueue,
public CompanySymbolService(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider) {
super(kafkaQueue);
super(kafkaQueue, kafkaProducer);
this.companySymbolsMap = imdgProvider.getImdg(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class);
}

View file

@ -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<Contact> contactMap;
@Autowired
public ContactService(Consumer<String, Object> kafkaQueue,
public ContactService(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider) {
super(kafkaQueue);
super(kafkaQueue, kafkaProducer);
this.contactMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Contact, Contact.class);
}

View file

@ -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
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

View file

@ -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<String, Object> createProducer(SecuritiesServiceSettings settings) {
return KafkaConsumerFactory.consumer(settings.getKafka());
public Consumer<String, Object> createConsumer(SecuritiesServiceSettings settings) {
return KafkaConsumerFactory.consumer(settings.getKafkaConsumer());
}
@Autowired
@Bean
public Producer<String, Object> createProducer(SecuritiesServiceSettings settings) {
return KafkaProducerFactory.producer(settings.getKafkaProducer());
}
}

View file

@ -4,6 +4,7 @@ import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.PropertySource;
import org.springframework.stereotype.Component;
import ru.spcex.clearing.platform.messaging.config.element.KafkaConsumerSettings;
import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings;
import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams;
@Component
@ -11,7 +12,8 @@ import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams;
@ConfigurationProperties("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;
}
}

View file

@ -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<MoneyMarketSecurity> moneyMarketSecurityMap;
private final ImdgProvider imdgProvider;
@Autowired
public MoneyMarketSecurityService(Consumer<String, Object> kafkaQueue, ImdgProvider imdgProvider) {
super(kafkaQueue);
public MoneyMarketSecurityService(Consumer<String, Object> kafkaQueue, Producer<String, Object> 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<RequestInfo> requestSpecificImdg = imdgProvider.getImdg(destination, RequestInfo.class);
Imdg<RequestInfo> 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<MoneyMarketSecurityNewRequest> userRequest) {
Imdg<RequestInfo> requestInfoImdg = imdgProvider.getImdg(Consts.DESTINATION_MONEY_MARKET_SECURITY_NEW, RequestInfo.class);
RequestInfo reqInfo = requestInfoImdg.getSingleObjectByID(userRequest.getId());
// Imdg<RequestInfo> 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<MoneyMarketSecurity> moneyMarketSecurityMap = transaction.getImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, MoneyMarketSecurity.class);
Imdg<RequestInfo> reqInfoMap = transaction.getImdg(Consts.DESTINATION_MONEY_MARKET_SECURITY_NEW, RequestInfo.class);
Imdg<RequestInfo> 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());
}

View file

@ -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());

View file

@ -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<String, Object> createProducer(UtilityServiceSettings settings) {
return KafkaConsumerFactory.consumer(settings.getKafka());
public Consumer<String, Object> createConsumer(UtilityServiceSettings settings) {
return KafkaConsumerFactory.consumer(settings.getKafkaConsumer());
}
@Autowired
@Bean
public Producer<String, Object> createProducer(UtilityServiceSettings settings) {
return KafkaProducerFactory.producer(settings.getKafkaProducer());
}
}

View file

@ -4,6 +4,7 @@ import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.PropertySource;
import org.springframework.stereotype.Component;
import ru.spcex.clearing.platform.messaging.config.element.KafkaConsumerSettings;
import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings;
import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams;
@Component
@ -11,7 +12,8 @@ import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams;
@ConfigurationProperties("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;
}
}

View file

@ -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<KeyRate> keyRateMap;
@Autowired
public KeyRateService(Consumer<String, Object> kafkaQueue, ImdgProvider imdgProvider) {
super(kafkaQueue);
public KeyRateService(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider) {
super(kafkaQueue, kafkaProducer);
this.keyRateMap = imdgProvider.getImdg(IMDGDistributedNames.Map_KeyRate, KeyRate.class);
}

View file

@ -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<UserSettings> userSettingsMap;
@Autowired
public UserSettingService(Consumer<String, Object> kafkaQueue, ImdgProvider imdgProvider) {
super(kafkaQueue);
public UserSettingService(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider) {
super(kafkaQueue, kafkaProducer);
this.userSettingsMap = imdgProvider.getImdg(IMDGDistributedNames.Map_UserSettings, UserSettings.class);
}

View file

@ -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
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

View file

@ -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<RecordMetadata> send;
//default response
BaseRequest<RequestInfoUpdate> 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<RecordMetadata> send;
BaseRequest<Object> 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));