From 30ec2ceeec4bf3cf8da5de1bc13e4bb99698a474 Mon Sep 17 00:00:00 2001 From: ialbert Date: Fri, 25 Nov 2022 12:58:10 +0300 Subject: [PATCH 1/5] http://jira.mfd.msk:8088/browse/CLS-51 --- scheduler service (addition of scheduler) --- clearing-parent/pom.xml | 1 + clearing-parent/scheduler-service/pom.xml | 81 ++++++++++++++++ .../SchedulerServiceApplication.java | 12 +++ .../scheduler/config/KafkaConfig.java | 29 ++++++ .../config/SchedulerServiceImdgConfig.java | 47 +++++++++ .../settings/SchedulerServiceSettings.java | 41 ++++++++ .../service/SchedulerAllTodayService.java | 82 ++++++++++++++++ .../scheduler/service/SchedulerService.java | 79 +++++++++++++++ .../platform/messaging/domain/Consts.java | 4 + .../cud/schedule/SchedulerNewRequest.java | 91 ++++++++++++++++++ .../cud/schedule/SchedulerUpdateRequest.java | 95 +++++++++++++++++++ 11 files changed, 562 insertions(+) create mode 100644 clearing-parent/scheduler-service/pom.xml create mode 100644 clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/SchedulerServiceApplication.java create mode 100644 clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/KafkaConfig.java create mode 100644 clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/SchedulerServiceImdgConfig.java create mode 100644 clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/settings/SchedulerServiceSettings.java create mode 100644 clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/SchedulerAllTodayService.java create mode 100644 clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/SchedulerService.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/SchedulerNewRequest.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/SchedulerUpdateRequest.java diff --git a/clearing-parent/pom.xml b/clearing-parent/pom.xml index 2ef81dddd..08cc8557b 100644 --- a/clearing-parent/pom.xml +++ b/clearing-parent/pom.xml @@ -29,6 +29,7 @@ reports-service account-service balance-service + scheduler-service diff --git a/clearing-parent/scheduler-service/pom.xml b/clearing-parent/scheduler-service/pom.xml new file mode 100644 index 000000000..04774fd14 --- /dev/null +++ b/clearing-parent/scheduler-service/pom.xml @@ -0,0 +1,81 @@ + + + + clearing-parent + ru.spcex.clearing + SPCEX-1.0.0.0 + + 4.0.0 + + scheduler-service + + + 17 + 17 + + + + ru.spcex.platform + platform-messaging + + + ru.spcex.platform + platform-enum + + + ru.spcex.platform + platform-imdg-api-hazelcast-impl + + + ru.spcex.clearing + classes + + + org.springframework.boot + spring-boot-starter + + + com.fasterxml.jackson.core + jackson-databind + + + + + org.springframework.boot + spring-boot-starter-test + test + + + + + + + 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/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/SchedulerServiceApplication.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/SchedulerServiceApplication.java new file mode 100644 index 000000000..6022af291 --- /dev/null +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/SchedulerServiceApplication.java @@ -0,0 +1,12 @@ +package ru.spcex.clearing.scheduler; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +@SpringBootApplication +public class SchedulerServiceApplication { + public static void main(String[] args) { + SpringApplication springApplication = new SpringApplication(SchedulerServiceApplication.class); + springApplication.run(args); + } +} diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/KafkaConfig.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/KafkaConfig.java new file mode 100644 index 000000000..5babe0b15 --- /dev/null +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/KafkaConfig.java @@ -0,0 +1,29 @@ +package ru.spcex.clearing.scheduler.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.scheduler.config.settings.SchedulerServiceSettings; + +@Configuration +public class KafkaConfig { + @Autowired + @Bean + @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) + public Consumer createConsumer(SchedulerServiceSettings settings) { + return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); + } + + @Autowired + @Bean + public Producer createProducer(SchedulerServiceSettings settings) { + return KafkaProducerFactory.producer(settings.getKafkaProducer()); + } + +} diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/SchedulerServiceImdgConfig.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/SchedulerServiceImdgConfig.java new file mode 100644 index 000000000..f9d5da568 --- /dev/null +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/SchedulerServiceImdgConfig.java @@ -0,0 +1,47 @@ +package ru.spcex.clearing.scheduler.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.scheduler.config.settings.SchedulerServiceSettings; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; + +@Configuration +public class SchedulerServiceImdgConfig { + 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, + SchedulerServiceSettings settings + ) { + return new HazelcastService(taskExecutorHazelcastClientInitializer, + taskExecutorIdGeneratorAwaiter, + settings.getHazelcast()); + } + +} diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/settings/SchedulerServiceSettings.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/settings/SchedulerServiceSettings.java new file mode 100644 index 000000000..8135f4e2d --- /dev/null +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/settings/SchedulerServiceSettings.java @@ -0,0 +1,41 @@ +package ru.spcex.clearing.scheduler.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.clearing.platform.messaging.config.element.KafkaProducerSettings; +import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams; + +@Component +@PropertySource("file:${spring.config.location}/application.properties") +@ConfigurationProperties("scheduler-service") +public class SchedulerServiceSettings { + private HazelcastClientParams hazelcast; + private KafkaConsumerSettings kafkaConsumer; + private KafkaProducerSettings kafkaProducer; + + public HazelcastClientParams getHazelcast() { + return hazelcast; + } + + public void setHazelcast(HazelcastClientParams hazelcast) { + this.hazelcast = hazelcast; + } + + public KafkaConsumerSettings getKafkaConsumer() { + return kafkaConsumer; + } + + 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/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/SchedulerAllTodayService.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/SchedulerAllTodayService.java new file mode 100644 index 000000000..faff5d228 --- /dev/null +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/SchedulerAllTodayService.java @@ -0,0 +1,82 @@ +package ru.spcex.clearing.scheduler.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; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.scheduler.Scheduler; +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.domain.cud.schedule.SchedulerNewRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.schedule.SchedulerUpdateRequest; +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 SchedulerAllTodayService extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg schedulerMap; + + @Autowired + public SchedulerAllTodayService(Consumer kafkaQueue, Producer kafkaProducer, + ImdgProvider imdgProvider) { + super(kafkaQueue, kafkaProducer); + this.schedulerMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Scheduler, Scheduler.class); + } + + @Override + public void afterPropertiesSet() { + callback(SchedulerNewRequest.class) + .setConsumer(this::newScheduler) + .forDestination(Consts.DESTINATION_SCHEDULER_NEW, callbacks::put); + callback(SchedulerUpdateRequest.class) + .setConsumer(this::updateScheduler) + .forDestination(Consts.DESTINATION_SCHEDULER_UPDATE, callbacks::put); + callback(CommonDeleteRequest.class) + .setConsumer(this::deleteScheduler) + .forDestination(Consts.DESTINATION_SCHEDULER_DELETE, callbacks::put); + init(); + } + + private void newScheduler(BaseRequest userRequest) { + SchedulerNewRequest req = userRequest.getRequestPayload(); + log.debug("SchedulerNewRequest received"); + Scheduler scheduler = new Scheduler(); + scheduler.setTask(req.getTask()); + scheduler.setTaskTime(req.getTaskTime()); + scheduler.setClearingDate(req.getClearingDate()); + scheduler.setMarket(req.getMarket()); + scheduler.setTaskStatus(req.getTaskStatus()); + scheduler.setSecurityId(req.getSecurityId()); + schedulerMap.insert(scheduler); + log.debug("successfully processed, new id {}", scheduler.getId()); + } + + private void updateScheduler(BaseRequest userRequest) { + SchedulerUpdateRequest req = userRequest.getRequestPayload(); + log.debug("SchedulerUpdateRequest received id = {}", req.getId()); + Scheduler scheduler = schedulerMap.getSingleObjectByID(req.getId()); + scheduler.setTask(req.getTask()); + scheduler.setTaskTime(req.getTaskTime()); + scheduler.setClearingDate(req.getClearingDate()); + scheduler.setMarket(req.getMarket()); + scheduler.setTaskStatus(req.getTaskStatus()); + scheduler.setSecurityId(req.getSecurityId()); + + schedulerMap.update(scheduler); + } + + private void deleteScheduler(BaseRequest userRequest) { + CommonDeleteRequest req = userRequest.getRequestPayload(); + log.debug("CommonDeleteRequest received id = {}", req.getId()); + Scheduler scheduler = schedulerMap.getSingleObjectByID(req.getId()); + schedulerMap.delete(scheduler); + } + +} diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/SchedulerService.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/SchedulerService.java new file mode 100644 index 000000000..effaaf686 --- /dev/null +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/SchedulerService.java @@ -0,0 +1,79 @@ +package ru.spcex.clearing.scheduler.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; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.misc.KeyRate; +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.domain.cud.utilities.KeyRateNewRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.utilities.KeyRateUpdateRequest; +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.platform.enumeration.Status; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +@Service +public class SchedulerService extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg keyRateMap; + + @Autowired + public SchedulerService(Consumer kafkaQueue, Producer kafkaProducer, + ImdgProvider imdgProvider) { + super(kafkaQueue, kafkaProducer); + this.keyRateMap = imdgProvider.getImdg(IMDGDistributedNames.Map_KeyRate, KeyRate.class); + } + + @Override + public void afterPropertiesSet() { + callback(KeyRateNewRequest.class) + .setConsumer(this::newKeyRate) + .forDestination(Consts.DESTINATION_KEY_RATE_NEW, callbacks::put); + callback(KeyRateUpdateRequest.class) + .setConsumer(this::updateKeyRate) + .forDestination(Consts.DESTINATION_KEY_RATE_UPDATE, callbacks::put); + callback(CommonDeleteRequest.class) + .setConsumer(this::deleteKeyRate) + .forDestination(Consts.DESTINATION_KEY_RATE_DELETE, callbacks::put); + init(); + } + + private void newKeyRate(BaseRequest userRequest) { + KeyRateNewRequest req = userRequest.getRequestPayload(); + log.debug("MoneyMarketSecurityNewRequest received"); + KeyRate keyRate = new KeyRate(); + keyRate.setEndDate(req.getEndDate()); + keyRate.setDocument(req.getDocument()); + keyRate.setRate(req.getKeyRate()); + keyRate.setStartDate(req.getStartDate()); + keyRate.setWorkflowStatus(Status.Active.getKey()); + keyRateMap.insert(keyRate); + log.debug("successfully processed, new id {}", keyRate.getId()); + } + + private void updateKeyRate(BaseRequest userRequest) { + KeyRateUpdateRequest req = userRequest.getRequestPayload(); + log.debug("MoneyMarketSecurityUpdateRequest received id = {}", req.getId()); + KeyRate keyRate = keyRateMap.getSingleObjectByID(req.getId()); + keyRate.setEndDate(req.getEndDate()); + keyRate.setDocument(req.getDocument()); + keyRate.setRate(req.getKeyRate()); + keyRate.setStartDate(req.getStartDate()); + keyRateMap.update(keyRate); + } + + private void deleteKeyRate(BaseRequest userRequest) { + CommonDeleteRequest req = userRequest.getRequestPayload(); + log.debug("CommonDeleteRequest received id = {}", req.getId()); + KeyRate keyRate = keyRateMap.getSingleObjectByID(req.getId()); + keyRateMap.delete(keyRate); + } + +} 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 759cb6328..df5ac80b5 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 @@ -9,6 +9,10 @@ public interface Consts { String DESTINATION_KEY_RATE_UPDATE = "key-rate-update"; String DESTINATION_KEY_RATE_DELETE = "key-rate-delete"; + String DESTINATION_SCHEDULER_NEW = "scheduler-new"; + String DESTINATION_SCHEDULER_UPDATE = "scheduler-update"; + String DESTINATION_SCHEDULER_DELETE = "scheduler-delete"; + String DESTINATION_COMPANY_DELETE = "company-delete"; String DESTINATION_COMPANY_INFO_UPDATE = "company-info-update"; String DESTINATION_COMPANY_SYMBOL_UPDATE = "company-symbol-update"; diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/SchedulerNewRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/SchedulerNewRequest.java new file mode 100644 index 000000000..405d98f59 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/SchedulerNewRequest.java @@ -0,0 +1,91 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.schedule; + +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.databind.annotation.JsonDeserialize; +import com.fasterxml.jackson.databind.annotation.JsonSerialize; +import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalDateDeserializer; +import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalDateSerializer; + +import java.time.LocalDate; +import java.time.LocalTime; + +public class SchedulerNewRequest { + @JsonProperty + @JsonSerialize(using = LocalDateSerializer.class) + @JsonDeserialize(using = LocalDateDeserializer.class) + public LocalTime taskTime; + + @JsonProperty + public String task; + + @JsonProperty + @JsonSerialize(using = LocalDateSerializer.class) + @JsonDeserialize(using = LocalDateDeserializer.class) + public LocalDate clearingDate; + + @JsonProperty + public String market; + + @JsonProperty + public String taskStatus; + + @JsonProperty + public Long securityId; + + public SchedulerNewRequest(LocalTime taskTime, String task, LocalDate clearingDate, String market, String taskStatus, Long securityId) { + this.taskTime = taskTime; + this.task = task; + this.clearingDate = clearingDate; + this.market = market; + this.taskStatus = taskStatus; + this.securityId = securityId; + } + + public LocalTime getTaskTime() { + return taskTime; + } + + public void setTaskTime(LocalTime taskTime) { + this.taskTime = taskTime; + } + + public String getTask() { + return task; + } + + public void setTask(String task) { + this.task = task; + } + + public LocalDate getClearingDate() { + return clearingDate; + } + + public void setClearingDate(LocalDate clearingDate) { + this.clearingDate = clearingDate; + } + + public String getMarket() { + return market; + } + + public void setMarket(String market) { + this.market = market; + } + + public String getTaskStatus() { + return taskStatus; + } + + public void setTaskStatus(String taskStatus) { + this.taskStatus = taskStatus; + } + + public Long getSecurityId() { + return securityId; + } + + public void setSecurityId(Long securityId) { + this.securityId = securityId; + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/SchedulerUpdateRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/SchedulerUpdateRequest.java new file mode 100644 index 000000000..abdf25d5a --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/SchedulerUpdateRequest.java @@ -0,0 +1,95 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.schedule; + +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.databind.annotation.JsonDeserialize; +import com.fasterxml.jackson.databind.annotation.JsonSerialize; +import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalDateDeserializer; +import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalTimeDeserializer; +import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalDateSerializer; +import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalTimeSerializer; + +import java.time.LocalDate; +import java.time.LocalTime; + +public class SchedulerUpdateRequest { + @JsonProperty + private String task; + + @JsonSerialize(using = LocalTimeSerializer.class) + @JsonDeserialize(using = LocalTimeDeserializer.class) + @JsonProperty + private LocalTime taskTime; + + @JsonSerialize(using = LocalDateSerializer.class) + @JsonDeserialize(using = LocalDateDeserializer.class) + @JsonProperty + private LocalDate clearingDate; + + @JsonProperty + private String market; + + @JsonProperty + private String taskStatus; + + @JsonProperty + private Long securityId; + + @JsonProperty + private Long id; + + public String getTask() { + return task; + } + + public void setTask(String task) { + this.task = task; + } + + public LocalTime getTaskTime() { + return taskTime; + } + + public void setTaskTime(LocalTime taskTime) { + this.taskTime = taskTime; + } + + public LocalDate getClearingDate() { + return clearingDate; + } + + public void setClearingDate(LocalDate clearingDate) { + this.clearingDate = clearingDate; + } + + public String getMarket() { + return market; + } + + public void setMarket(String market) { + this.market = market; + } + + public String getTaskStatus() { + return taskStatus; + } + + public void setTaskStatus(String taskStatus) { + this.taskStatus = taskStatus; + } + + public Long getSecurityId() { + return securityId; + } + + public void setSecurityId(Long securityId) { + this.securityId = securityId; + } + + public Long getId() { + return id; + } + + public void setId(Long id) { + this.id = id; + } +} From d7d098990b839910a6c00811861867cda085d610 Mon Sep 17 00:00:00 2001 From: aalehin Date: Tue, 29 Nov 2022 10:56:23 +0300 Subject: [PATCH 2/5] http://jira.mfd.msk:8088/browse/CLS-51 --- 1) application.properties 2) logback.xml --- .../src/main/resources/application.properties | 17 +++++++++ .../src/main/resources/logback.xml | 38 +++++++++++++++++++ 2 files changed, 55 insertions(+) create mode 100644 clearing-parent/scheduler-service/src/main/resources/application.properties create mode 100644 clearing-parent/scheduler-service/src/main/resources/logback.xml diff --git a/clearing-parent/scheduler-service/src/main/resources/application.properties b/clearing-parent/scheduler-service/src/main/resources/application.properties new file mode 100644 index 000000000..744aa5f7d --- /dev/null +++ b/clearing-parent/scheduler-service/src/main/resources/application.properties @@ -0,0 +1,17 @@ +spring.main.web-application-type=none +scheduler-service.hazelcast.cluster-members=127.0.0.1:5701 +scheduler-service.hazelcast.login=dev +scheduler-service.hazelcast.password=dev-pass +scheduler-service.kafka-consumer.bootstrap-servers=localhost:9092 +scheduler-service.kafka-consumer.group-id=dev-group-scheduler-service +scheduler-service.kafka-consumer.enable-auto-commit=true +scheduler-service.kafka-consumer.session-timeout-ms=30000 +scheduler-service.kafka-consumer.auto-offset-reset=latest +scheduler-service.kafka-consumer.linger-ms=1 +scheduler-service.kafka-consumer.buffer-memory=33554432 +scheduler-service.kafka-producer.bootstrap-servers=localhost:9092 +scheduler-service.kafka-producer.acks=all +scheduler-service.kafka-producer.retries=0 +scheduler-service.kafka-producer.batch-size=16384 +scheduler-service.kafka-producer.linger-ms=1 +scheduler-service.kafka-producer.buffer-memory=33554432 \ No newline at end of file diff --git a/clearing-parent/scheduler-service/src/main/resources/logback.xml b/clearing-parent/scheduler-service/src/main/resources/logback.xml new file mode 100644 index 000000000..a7a54b222 --- /dev/null +++ b/clearing-parent/scheduler-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/utility-service.log + + + %d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %class{0}:%msg%n + utf8 + + + + ./logs/utility-service.%i.log + + 1 + 10 + + + 500MB + + + + + + + + + + + + + From 86e9b724ac3832d090dca51471e93646893e18c3 Mon Sep 17 00:00:00 2001 From: aalehin Date: Tue, 29 Nov 2022 16:36:12 +0300 Subject: [PATCH 3/5] http://jira.mfd.msk:8088/browse/CLS-51 --- SchedulerAllTodayService fix up --- .../service/SchedulerAllTodayService.java | 55 +------------------ 1 file changed, 3 insertions(+), 52 deletions(-) diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/SchedulerAllTodayService.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/SchedulerAllTodayService.java index faff5d228..959c40453 100644 --- a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/SchedulerAllTodayService.java +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/SchedulerAllTodayService.java @@ -7,13 +7,8 @@ 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.scheduler.Scheduler; +import ru.clearing.classes.statics.data.scheduler.SchedulerAllToday; 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.domain.cud.schedule.SchedulerNewRequest; -import ru.spcex.clearing.platform.messaging.domain.cud.schedule.SchedulerUpdateRequest; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; @@ -21,62 +16,18 @@ import ru.spcex.platform.imdg.api.ImdgProvider; @Service public class SchedulerAllTodayService extends QueueConsumer implements InitializingBean { private final Logger log = LoggerFactory.getLogger(getClass()); - private final Imdg schedulerMap; + private final Imdg schedulerAllTodayMap; @Autowired public SchedulerAllTodayService(Consumer kafkaQueue, Producer kafkaProducer, ImdgProvider imdgProvider) { super(kafkaQueue, kafkaProducer); - this.schedulerMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Scheduler, Scheduler.class); + this.schedulerAllTodayMap = imdgProvider.getImdg(IMDGDistributedNames.Map_SchedulerAllToday, SchedulerAllToday.class); } @Override public void afterPropertiesSet() { - callback(SchedulerNewRequest.class) - .setConsumer(this::newScheduler) - .forDestination(Consts.DESTINATION_SCHEDULER_NEW, callbacks::put); - callback(SchedulerUpdateRequest.class) - .setConsumer(this::updateScheduler) - .forDestination(Consts.DESTINATION_SCHEDULER_UPDATE, callbacks::put); - callback(CommonDeleteRequest.class) - .setConsumer(this::deleteScheduler) - .forDestination(Consts.DESTINATION_SCHEDULER_DELETE, callbacks::put); init(); } - private void newScheduler(BaseRequest userRequest) { - SchedulerNewRequest req = userRequest.getRequestPayload(); - log.debug("SchedulerNewRequest received"); - Scheduler scheduler = new Scheduler(); - scheduler.setTask(req.getTask()); - scheduler.setTaskTime(req.getTaskTime()); - scheduler.setClearingDate(req.getClearingDate()); - scheduler.setMarket(req.getMarket()); - scheduler.setTaskStatus(req.getTaskStatus()); - scheduler.setSecurityId(req.getSecurityId()); - schedulerMap.insert(scheduler); - log.debug("successfully processed, new id {}", scheduler.getId()); - } - - private void updateScheduler(BaseRequest userRequest) { - SchedulerUpdateRequest req = userRequest.getRequestPayload(); - log.debug("SchedulerUpdateRequest received id = {}", req.getId()); - Scheduler scheduler = schedulerMap.getSingleObjectByID(req.getId()); - scheduler.setTask(req.getTask()); - scheduler.setTaskTime(req.getTaskTime()); - scheduler.setClearingDate(req.getClearingDate()); - scheduler.setMarket(req.getMarket()); - scheduler.setTaskStatus(req.getTaskStatus()); - scheduler.setSecurityId(req.getSecurityId()); - - schedulerMap.update(scheduler); - } - - private void deleteScheduler(BaseRequest userRequest) { - CommonDeleteRequest req = userRequest.getRequestPayload(); - log.debug("CommonDeleteRequest received id = {}", req.getId()); - Scheduler scheduler = schedulerMap.getSingleObjectByID(req.getId()); - schedulerMap.delete(scheduler); - } - } From e8bc16553774bd4d7792f91fd7a989d17590c689 Mon Sep 17 00:00:00 2001 From: ialbert Date: Tue, 29 Nov 2022 16:50:55 +0300 Subject: [PATCH 4/5] http://jira.mfd.msk:8088/browse/CLS-97 --- .../java/ru/clearing/classes/statics/data/sdf/SDf03.java | 9 +++++++++ .../ru/spcex/clearing/imdg/object/SDf03MapStore.java | 4 +++- 2 files changed, 12 insertions(+), 1 deletion(-) diff --git a/clearing-parent/classes/src/main/java/ru/clearing/classes/statics/data/sdf/SDf03.java b/clearing-parent/classes/src/main/java/ru/clearing/classes/statics/data/sdf/SDf03.java index 566690b89..4d190fef9 100644 --- a/clearing-parent/classes/src/main/java/ru/clearing/classes/statics/data/sdf/SDf03.java +++ b/clearing-parent/classes/src/main/java/ru/clearing/classes/statics/data/sdf/SDf03.java @@ -16,6 +16,7 @@ public class SDf03 extends SpcexObjectBase { private String seg_type; private String doc_type; private String docnm_ref; + private String docnmprev; private String priority; private String sbankcode; private String c_acc_deb; @@ -436,4 +437,12 @@ public class SDf03 extends SpcexObjectBase { public void setGenerationId(Long generationId) { this.generationId = generationId; } + + public String getDocnmprev() { + return docnmprev; + } + + public void setDocnmprev(String docnmprev) { + this.docnmprev = docnmprev; + } } diff --git a/clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/object/SDf03MapStore.java b/clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/object/SDf03MapStore.java index f549d9bff..d7b3fe94e 100644 --- a/clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/object/SDf03MapStore.java +++ b/clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/object/SDf03MapStore.java @@ -30,7 +30,7 @@ public class SDf03MapStore extends TemplateMapStore { @Override public String[] getFields() { return new String[] { - "ID", "SEG_TYPE", "DOC_TYPE", "DOCNM_REF", "PRIORITY", "SBANKCODE", "C_ACC_DEB", "SBANKNAM1", "SBANKNAM2", + "ID", "SEG_TYPE", "DOC_TYPE", "DOCNM_REF", "DOCNMPREV", "PRIORITY", "SBANKCODE", "C_ACC_DEB", "SBANKNAM1", "SBANKNAM2", "SBANKNAM3", "SBANKNAM4", "SBANKNAM5", "RBANKCODE", "C_ACC_CRED", "RBANKNAM1", "RBANKNAM2", "RBANKNAM3", "RBANKNAM4", "RBANKNAM5", "PAY_DATE", "EXT_DATE", "PAY_VAL", "SUM_DEB", "SCLIENTN1", "SCLIENTN2", "SCLIENTN3", "SCLIENTN4", "SC_CODE", "ACC_DEB", "RCLIENTN1", "RCLIENTN2", "RCLIENTN3", "RCLIENTN4", @@ -46,6 +46,7 @@ public class SDf03MapStore extends TemplateMapStore { object.setSeg_type(resultSet.getObject("SEG_TYPE", String.class)); object.setDoc_type(resultSet.getObject("DOC_TYPE", String.class)); object.setDocnm_ref(resultSet.getObject("DOCNM_REF", String.class)); + object.setDocnmprev(resultSet.getObject("DOCNMPREV", String.class)); object.setPriority(resultSet.getObject("PRIORITY", String.class)); object.setSbankcode(resultSet.getObject("SBANKCODE", String.class)); object.setC_acc_deb(resultSet.getObject("C_ACC_DEB", String.class)); @@ -100,6 +101,7 @@ public class SDf03MapStore extends TemplateMapStore { object.getSeg_type(), object.getDoc_type(), object.getDocnm_ref(), + object.getDocnmprev(), object.getPriority(), object.getSbankcode(), object.getC_acc_deb(), From 49729c4d0893c438a4a3202ac0c472d55f4088fd Mon Sep 17 00:00:00 2001 From: aalehin Date: Tue, 29 Nov 2022 19:10:53 +0300 Subject: [PATCH 5/5] http://jira.mfd.msk:8088/browse/CLS-51 --- some fix up --- .../AbstractHazelcastLifecycleSupport.java | 29 +++++++++--------- .../ru/spcex/clearing/imdg/util/Util.java | 30 +++++++++++++++++++ 2 files changed, 44 insertions(+), 15 deletions(-) create mode 100644 clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/util/Util.java diff --git a/clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/services/AbstractHazelcastLifecycleSupport.java b/clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/services/AbstractHazelcastLifecycleSupport.java index 5c1e24129..e55d066a8 100644 --- a/clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/services/AbstractHazelcastLifecycleSupport.java +++ b/clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/services/AbstractHazelcastLifecycleSupport.java @@ -19,30 +19,30 @@ import java.util.ArrayList; import java.util.List; import java.util.concurrent.*; +import static ru.spcex.clearing.imdg.util.Util.makeSchedulerAllTodayMap; + //todo почистить класс public abstract class AbstractHazelcastLifecycleSupport implements InitializingBean, DisposableBean { - /** - * Рабочая версия БД. Треьуется вручную сверять с DDL.sql и накручивать эту переменную. - */ - public abstract String getCheckDbVersion(); - - /** - * Обязательная сверка у этих классов SerialVersionUID во время подключения к Storage. - */ - public abstract Class[] getSerialVersionUIDClasses(); - private final Logger log = LoggerFactory.getLogger(this.getClass()); - private final HazelcastInstance hazelcastServerInstance; private final JdbcTemplate jdbcTemplate; -// private final HazelcastClientListener clientListener; -// private final ITaskAdministrator startupTasksAdministrator; public AbstractHazelcastLifecycleSupport(HazelcastInstance hazelcastServerInstance, JdbcTemplate jdbcTemplate) { this.hazelcastServerInstance = hazelcastServerInstance; this.jdbcTemplate = jdbcTemplate; } + /** + * Рабочая версия БД. Треьуется вручную сверять с DDL.sql и накручивать эту переменную. + */ + public abstract String getCheckDbVersion(); +// private final HazelcastClientListener clientListener; +// private final ITaskAdministrator startupTasksAdministrator; + + /** + * Обязательная сверка у этих классов SerialVersionUID во время подключения к Storage. + */ + public abstract Class[] getSerialVersionUIDClasses(); @Override public void afterPropertiesSet() { @@ -99,7 +99,7 @@ public abstract class AbstractHazelcastLifecycleSupport implements InitializingB boolean generatorResult = generator.init(maxKey); if (generatorResult) { log.info("IDGenerator {} success init by {}", IMDGDistributedNames.MAP_SEQUENCE_NAME, maxKey); -// makeSchedulerAllTodayMap(); + makeSchedulerAllTodayMap(hazelcastServerInstance); } else { log.info("IDGenerator {} already initialized in other node", IMDGDistributedNames.MAP_SEQUENCE_NAME); } @@ -114,7 +114,6 @@ public abstract class AbstractHazelcastLifecycleSupport implements InitializingB } - @Override public void destroy() { hazelcastServerInstance.shutdown(); diff --git a/clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/util/Util.java b/clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/util/Util.java new file mode 100644 index 000000000..432c51755 --- /dev/null +++ b/clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/util/Util.java @@ -0,0 +1,30 @@ +package ru.spcex.clearing.imdg.util; + +import com.hazelcast.core.HazelcastInstance; +import com.hazelcast.core.IMap; +import ru.clearing.classes.statics.data.scheduler.Scheduler; +import ru.clearing.classes.statics.data.scheduler.Timetable; +import ru.clearing.classes.statics.data.scheduler.TradingCalendar; + +import java.time.LocalDate; +import java.util.List; + +import static ru.spcex.clearing.imdg.IMDGDistributedNames.*; + +public class Util { + + public static void makeSchedulerAllTodayMap(HazelcastInstance hazelcastInstance) { + IMap schedulerMap = hazelcastInstance.getMap(Map_Scheduler); + IMap tradingCalendarMap = hazelcastInstance.getMap(Map_TradingCalendar); + IMap timetableMap = hazelcastInstance.getMap(Map_Timetable); + Boolean isTradingCalendarEmpty = tradingCalendarMap.isEmpty(); + + List listOfSchedulerOnDate = schedulerMap.values().stream().filter((x) -> x.getClearingDate().isEqual(LocalDate.now())).toList(); + if (isTradingCalendarEmpty) { + //то из timeTable все активные записи. + } else { + + } + + } +}