diff --git a/clearing-parent/company-service/pom.xml b/clearing-parent/company-service/pom.xml new file mode 100644 index 000000000..cd323c846 --- /dev/null +++ b/clearing-parent/company-service/pom.xml @@ -0,0 +1,70 @@ + + + + clearing-parent + ru.spcex.clearing + SPCEX-1.0.0.0 + + 4.0.0 + + company-service + + + 17 + 17 + + + + ru.spcex.platform + platform-messaging + + + ru.spcex.platform + platform-imdg-api-hazelcast-impl + + + ru.spcex.clearing + classes + + + org.springframework.boot + spring-boot-starter + + + com.fasterxml.jackson.core + jackson-databind + + + + + + + src/main/resources + + application.properties + + false + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + repackage + + + + + ${project.artifactId} + + + + + + + \ No newline at end of file diff --git a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/CompanyServiceApplication.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/CompanyServiceApplication.java new file mode 100644 index 000000000..19bdaf7b1 --- /dev/null +++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/CompanyServiceApplication.java @@ -0,0 +1,12 @@ +package ru.spcex.clearing.company; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +@SpringBootApplication +public class CompanyServiceApplication { + public static void main(String[] args) { + SpringApplication app = new SpringApplication(CompanyServiceApplication.class); + app.run(args); + } +} 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 new file mode 100644 index 000000000..d63ecf6c5 --- /dev/null +++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/KafkaConfig.java @@ -0,0 +1,17 @@ +package ru.spcex.clearing.company.config; + +import org.apache.kafka.clients.consumer.Consumer; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import ru.spcex.clearing.company.config.settings.CompanyServiceSettings; +import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; + +@Configuration +public class KafkaConfig { + @Autowired + @Bean + public Consumer createProducer(CompanyServiceSettings settings) { + return KafkaConsumerFactory.consumer(settings.getKafka()); + } +} diff --git a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/UtilityServiceImdgConfig.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/UtilityServiceImdgConfig.java new file mode 100644 index 000000000..bbe41d92f --- /dev/null +++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/UtilityServiceImdgConfig.java @@ -0,0 +1,47 @@ +package ru.spcex.clearing.company.config; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; +import ru.spcex.clearing.company.config.settings.CompanyServiceSettings; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; + +@Configuration +public class UtilityServiceImdgConfig { + private static ThreadPoolTaskExecutor createThreadPoolTaskExecutor(int maxPoolSz, boolean waitForCompletion) { + ThreadPoolTaskExecutor pool = new ThreadPoolTaskExecutor(); + if (maxPoolSz > 2) { + pool.setKeepAliveSeconds(60); + pool.setAllowCoreThreadTimeOut(true); + } + pool.setCorePoolSize(maxPoolSz); + pool.setWaitForTasksToCompleteOnShutdown(waitForCompletion); + return pool; + } + + @Bean(name = "taskExecutorHazelcastClientInitializer") + public ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer() { + return createThreadPoolTaskExecutor(1, true); + } + + @Bean(name = "taskExecutorIdGeneratorAwaiter") + public ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter() { + return createThreadPoolTaskExecutor(1, false); + } + + @Autowired + @Bean + public ImdgProvider imdgProvider( + @Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer, + @Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter, + CompanyServiceSettings settings + ) { + return new HazelcastService(taskExecutorHazelcastClientInitializer, + taskExecutorIdGeneratorAwaiter, + settings.getHazelcast()); + } + +} diff --git a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/service/CompanyInfoService.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/service/CompanyInfoService.java new file mode 100644 index 000000000..ae055a49b --- /dev/null +++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/service/CompanyInfoService.java @@ -0,0 +1,56 @@ +package ru.spcex.clearing.company.config.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.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.profile.CompanyInfo; +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.domain.cud.company.CompanyInfoUpdateRequest; +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +@Service +public class CompanyInfoService extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg companyInfoMap; + + @Autowired + public CompanyInfoService(Consumer kafkaQueue, ImdgProvider imdgProvider) { + super(kafkaQueue); + this.companyInfoMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, CompanyInfo.class); + } + + @Override + public void afterPropertiesSet() { + callback(CompanyInfoUpdateRequest.class) + .setConsumer(this::companyInfoUpdate) + .forDestination(Consts.DESTINATION_COMPANY_INFO_UPDATE, callbacks::put); + init(); + } + + public void companyInfoUpdate(BaseRequest userRequest) { + CompanyInfoUpdateRequest req = userRequest.getRequestPayload(); + log.debug("CompanyInfoUpdateRequest received"); + CompanyInfo companyInfo = companyInfoMap.getSingleObjectByID(req.getId()); + + companyInfo.setCorporationSoleType(req.getCorporationSoleType()); + companyInfo.setCountryCode(req.getCountryCode()); + companyInfo.setDescription(req.getDescription()); + companyInfo.setProfessionalSign(req.getProfessionalSign()); + companyInfo.setLegalKind(req.getLegalKind()); + companyInfo.setOrganizationType(req.getOrganizationType()); + companyInfo.setResidence(req.getResidence()); + companyInfo.setShortNameEng(req.getShortNameEng()); + companyInfo.setFullNameEng(req.getFullNameEng()); + companyInfo.setShortName(req.getShortName()); + companyInfo.setFullName(req.getFullName()); + companyInfoMap.update(companyInfo); + log.debug("successfully processed, id {}", companyInfo.getId()); + } +} diff --git a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/service/CompanyService.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/service/CompanyService.java new file mode 100644 index 000000000..f53fe0fdc --- /dev/null +++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/service/CompanyService.java @@ -0,0 +1,43 @@ +package ru.spcex.clearing.company.config.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.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.company.Company; +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.domain.cud.common.CommonDeleteRequest; +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +@Service +public class CompanyService extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg companyMap; + + @Autowired + public CompanyService(Consumer kafkaQueue, ImdgProvider imdgProvider) { + super(kafkaQueue); + this.companyMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class); + } + + @Override + public void afterPropertiesSet() { + callback(CommonDeleteRequest.class) + .setConsumer(this::deleteCompany) + .forDestination(Consts.DESTINATION_COMPANY_DELETE, callbacks::put); + init(); + } + + private void deleteCompany(BaseRequest userRequest) { + CommonDeleteRequest req = userRequest.getRequestPayload(); + log.debug("CommonDeleteRequest received id = {}", req.getId()); + Company company = companyMap.getSingleObjectByID(req.getId()); + companyMap.delete(company); + } +} diff --git a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/service/CompanySymbolService.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/service/CompanySymbolService.java new file mode 100644 index 000000000..808f78a6e --- /dev/null +++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/service/CompanySymbolService.java @@ -0,0 +1,48 @@ +package ru.spcex.clearing.company.config.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.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.company.CompanySymbols; +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.domain.cud.company.CompanySymbolUpdateRequest; +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +@Service +public class CompanySymbolService extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg companySymbolsMap; + + @Autowired + public CompanySymbolService(Consumer kafkaQueue, ImdgProvider imdgProvider) { + super(kafkaQueue); + this.companySymbolsMap = imdgProvider.getImdg(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class); + } + + @Override + public void afterPropertiesSet() { + callback(CompanySymbolUpdateRequest.class) + .setConsumer(this::companySymbolUpdate) + .forDestination(Consts.DESTINATION_COMPANY_SYMBOL_UPDATE, callbacks::put); + init(); + } + + public void companySymbolUpdate(BaseRequest userRequest) { + CompanySymbolUpdateRequest req = userRequest.getRequestPayload(); + log.debug("CompanySymbolUpdateRequest received"); + CompanySymbols companySymbols = companySymbolsMap.getSingleObjectByID(req.getId()); + companySymbols.setCompanySymbol(req.getCompanySymbol()); + companySymbols.setCompanySymbolValue(req.getCompanySymbolValue()); + + companySymbolsMap.update(companySymbols); + log.debug("successfully processed, id {}", companySymbols.getId()); + } +} 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 new file mode 100644 index 000000000..98f925d5f --- /dev/null +++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/settings/CompanyServiceSettings.java @@ -0,0 +1,31 @@ +package ru.spcex.clearing.company.config.settings; + +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.platform.imdg.iml.hazelcast.config.HazelcastClientParams; + +@Component +@PropertySource("file:${spring.config.location}/application.properties") +@ConfigurationProperties("utility-service") +public class CompanyServiceSettings { + private HazelcastClientParams hazelcast; + private KafkaConsumerSettings kafka; + + public HazelcastClientParams getHazelcast() { + return hazelcast; + } + + public void setHazelcast(HazelcastClientParams hazelcast) { + this.hazelcast = hazelcast; + } + + public KafkaConsumerSettings getKafka() { + return kafka; + } + + public void setKafka(KafkaConsumerSettings kafka) { + this.kafka = kafka; + } +} diff --git a/clearing-parent/company-service/src/main/resources/application.properties b/clearing-parent/company-service/src/main/resources/application.properties new file mode 100644 index 000000000..f3f6bd027 --- /dev/null +++ b/clearing-parent/company-service/src/main/resources/application.properties @@ -0,0 +1,13 @@ +spring.main.web-application-type=none + +utility-service.hazelcast.cluster-members=127.0.0.1 +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.kafka.enable-auto-commit=false +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 diff --git a/clearing-parent/company-service/src/main/resources/logback.xml b/clearing-parent/company-service/src/main/resources/logback.xml new file mode 100644 index 000000000..bc21772aa --- /dev/null +++ b/clearing-parent/company-service/src/main/resources/logback.xml @@ -0,0 +1,38 @@ + + + + + + %date{HH:mm:ss.SSS} [%thread] %-5level %class{0}:%line - %message%n + utf-8 + + + + ./logs/company-service.log + + + %d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %class{0}:%msg%n + utf8 + + + + ./logs/company-service.%i.log + + 1 + 10 + + + 500MB + + + + + + + + + + + + + diff --git a/clearing-parent/pom.xml b/clearing-parent/pom.xml index b0cdf7411..130d7fa45 100644 --- a/clearing-parent/pom.xml +++ b/clearing-parent/pom.xml @@ -25,6 +25,7 @@ dbf-exporter securities-service utility-service + company-service