Merge branch 'dev' into psemenkov
This commit is contained in:
commit
0920088faf
17 changed files with 624 additions and 16 deletions
|
|
@ -16,6 +16,7 @@ public class SDf03 extends SpcexObjectBase {
|
||||||
private String seg_type;
|
private String seg_type;
|
||||||
private String doc_type;
|
private String doc_type;
|
||||||
private String docnm_ref;
|
private String docnm_ref;
|
||||||
|
private String docnmprev;
|
||||||
private String priority;
|
private String priority;
|
||||||
private String sbankcode;
|
private String sbankcode;
|
||||||
private String c_acc_deb;
|
private String c_acc_deb;
|
||||||
|
|
@ -436,4 +437,12 @@ public class SDf03 extends SpcexObjectBase {
|
||||||
public void setGenerationId(Long generationId) {
|
public void setGenerationId(Long generationId) {
|
||||||
this.generationId = generationId;
|
this.generationId = generationId;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public String getDocnmprev() {
|
||||||
|
return docnmprev;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setDocnmprev(String docnmprev) {
|
||||||
|
this.docnmprev = docnmprev;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -30,7 +30,7 @@ public class SDf03MapStore extends TemplateMapStore<SDf03> {
|
||||||
@Override
|
@Override
|
||||||
public String[] getFields() {
|
public String[] getFields() {
|
||||||
return new String[] {
|
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",
|
"SBANKNAM3", "SBANKNAM4", "SBANKNAM5", "RBANKCODE", "C_ACC_CRED", "RBANKNAM1", "RBANKNAM2", "RBANKNAM3",
|
||||||
"RBANKNAM4", "RBANKNAM5", "PAY_DATE", "EXT_DATE", "PAY_VAL", "SUM_DEB", "SCLIENTN1", "SCLIENTN2",
|
"RBANKNAM4", "RBANKNAM5", "PAY_DATE", "EXT_DATE", "PAY_VAL", "SUM_DEB", "SCLIENTN1", "SCLIENTN2",
|
||||||
"SCLIENTN3", "SCLIENTN4", "SC_CODE", "ACC_DEB", "RCLIENTN1", "RCLIENTN2", "RCLIENTN3", "RCLIENTN4",
|
"SCLIENTN3", "SCLIENTN4", "SC_CODE", "ACC_DEB", "RCLIENTN1", "RCLIENTN2", "RCLIENTN3", "RCLIENTN4",
|
||||||
|
|
@ -46,6 +46,7 @@ public class SDf03MapStore extends TemplateMapStore<SDf03> {
|
||||||
object.setSeg_type(resultSet.getObject("SEG_TYPE", String.class));
|
object.setSeg_type(resultSet.getObject("SEG_TYPE", String.class));
|
||||||
object.setDoc_type(resultSet.getObject("DOC_TYPE", String.class));
|
object.setDoc_type(resultSet.getObject("DOC_TYPE", String.class));
|
||||||
object.setDocnm_ref(resultSet.getObject("DOCNM_REF", 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.setPriority(resultSet.getObject("PRIORITY", String.class));
|
||||||
object.setSbankcode(resultSet.getObject("SBANKCODE", String.class));
|
object.setSbankcode(resultSet.getObject("SBANKCODE", String.class));
|
||||||
object.setC_acc_deb(resultSet.getObject("C_ACC_DEB", String.class));
|
object.setC_acc_deb(resultSet.getObject("C_ACC_DEB", String.class));
|
||||||
|
|
@ -100,6 +101,7 @@ public class SDf03MapStore extends TemplateMapStore<SDf03> {
|
||||||
object.getSeg_type(),
|
object.getSeg_type(),
|
||||||
object.getDoc_type(),
|
object.getDoc_type(),
|
||||||
object.getDocnm_ref(),
|
object.getDocnm_ref(),
|
||||||
|
object.getDocnmprev(),
|
||||||
object.getPriority(),
|
object.getPriority(),
|
||||||
object.getSbankcode(),
|
object.getSbankcode(),
|
||||||
object.getC_acc_deb(),
|
object.getC_acc_deb(),
|
||||||
|
|
|
||||||
|
|
@ -19,30 +19,30 @@ import java.util.ArrayList;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.concurrent.*;
|
import java.util.concurrent.*;
|
||||||
|
|
||||||
|
import static ru.spcex.clearing.imdg.util.Util.makeSchedulerAllTodayMap;
|
||||||
|
|
||||||
//todo почистить класс
|
//todo почистить класс
|
||||||
public abstract class AbstractHazelcastLifecycleSupport implements InitializingBean, DisposableBean {
|
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 Logger log = LoggerFactory.getLogger(this.getClass());
|
||||||
|
|
||||||
private final HazelcastInstance hazelcastServerInstance;
|
private final HazelcastInstance hazelcastServerInstance;
|
||||||
private final JdbcTemplate jdbcTemplate;
|
private final JdbcTemplate jdbcTemplate;
|
||||||
// private final HazelcastClientListener clientListener;
|
|
||||||
// private final ITaskAdministrator startupTasksAdministrator;
|
|
||||||
|
|
||||||
public AbstractHazelcastLifecycleSupport(HazelcastInstance hazelcastServerInstance, JdbcTemplate jdbcTemplate) {
|
public AbstractHazelcastLifecycleSupport(HazelcastInstance hazelcastServerInstance, JdbcTemplate jdbcTemplate) {
|
||||||
this.hazelcastServerInstance = hazelcastServerInstance;
|
this.hazelcastServerInstance = hazelcastServerInstance;
|
||||||
this.jdbcTemplate = jdbcTemplate;
|
this.jdbcTemplate = jdbcTemplate;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Рабочая версия БД. Треьуется вручную сверять с DDL.sql и накручивать эту переменную.
|
||||||
|
*/
|
||||||
|
public abstract String getCheckDbVersion();
|
||||||
|
// private final HazelcastClientListener clientListener;
|
||||||
|
// private final ITaskAdministrator startupTasksAdministrator;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Обязательная сверка у этих классов SerialVersionUID во время подключения к Storage.
|
||||||
|
*/
|
||||||
|
public abstract Class[] getSerialVersionUIDClasses();
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public void afterPropertiesSet() {
|
public void afterPropertiesSet() {
|
||||||
|
|
@ -99,7 +99,7 @@ public abstract class AbstractHazelcastLifecycleSupport implements InitializingB
|
||||||
boolean generatorResult = generator.init(maxKey);
|
boolean generatorResult = generator.init(maxKey);
|
||||||
if (generatorResult) {
|
if (generatorResult) {
|
||||||
log.info("IDGenerator {} success init by {}", IMDGDistributedNames.MAP_SEQUENCE_NAME, maxKey);
|
log.info("IDGenerator {} success init by {}", IMDGDistributedNames.MAP_SEQUENCE_NAME, maxKey);
|
||||||
// makeSchedulerAllTodayMap();
|
makeSchedulerAllTodayMap(hazelcastServerInstance);
|
||||||
} else {
|
} else {
|
||||||
log.info("IDGenerator {} already initialized in other node", IMDGDistributedNames.MAP_SEQUENCE_NAME);
|
log.info("IDGenerator {} already initialized in other node", IMDGDistributedNames.MAP_SEQUENCE_NAME);
|
||||||
}
|
}
|
||||||
|
|
@ -114,7 +114,6 @@ public abstract class AbstractHazelcastLifecycleSupport implements InitializingB
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public void destroy() {
|
public void destroy() {
|
||||||
hazelcastServerInstance.shutdown();
|
hazelcastServerInstance.shutdown();
|
||||||
|
|
|
||||||
|
|
@ -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<Long, Scheduler> schedulerMap = hazelcastInstance.getMap(Map_Scheduler);
|
||||||
|
IMap<Long, TradingCalendar> tradingCalendarMap = hazelcastInstance.getMap(Map_TradingCalendar);
|
||||||
|
IMap<Long, Timetable> timetableMap = hazelcastInstance.getMap(Map_Timetable);
|
||||||
|
Boolean isTradingCalendarEmpty = tradingCalendarMap.isEmpty();
|
||||||
|
|
||||||
|
List<Scheduler> listOfSchedulerOnDate = schedulerMap.values().stream().filter((x) -> x.getClearingDate().isEqual(LocalDate.now())).toList();
|
||||||
|
if (isTradingCalendarEmpty) {
|
||||||
|
//то из timeTable все активные записи.
|
||||||
|
} else {
|
||||||
|
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -29,6 +29,7 @@
|
||||||
<module>reports-service</module>
|
<module>reports-service</module>
|
||||||
<module>account-service</module>
|
<module>account-service</module>
|
||||||
<module>balance-service</module>
|
<module>balance-service</module>
|
||||||
|
<module>scheduler-service</module>
|
||||||
</modules>
|
</modules>
|
||||||
|
|
||||||
<properties>
|
<properties>
|
||||||
|
|
|
||||||
81
clearing-parent/scheduler-service/pom.xml
Normal file
81
clearing-parent/scheduler-service/pom.xml
Normal file
|
|
@ -0,0 +1,81 @@
|
||||||
|
<?xml version="1.0" encoding="UTF-8"?>
|
||||||
|
<project xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
|
||||||
|
xmlns="http://maven.apache.org/POM/4.0.0"
|
||||||
|
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
|
||||||
|
<parent>
|
||||||
|
<artifactId>clearing-parent</artifactId>
|
||||||
|
<groupId>ru.spcex.clearing</groupId>
|
||||||
|
<version>SPCEX-1.0.0.0</version>
|
||||||
|
</parent>
|
||||||
|
<modelVersion>4.0.0</modelVersion>
|
||||||
|
|
||||||
|
<artifactId>scheduler-service</artifactId>
|
||||||
|
|
||||||
|
<properties>
|
||||||
|
<maven.compiler.source>17</maven.compiler.source>
|
||||||
|
<maven.compiler.target>17</maven.compiler.target>
|
||||||
|
</properties>
|
||||||
|
<dependencies>
|
||||||
|
<dependency>
|
||||||
|
<groupId>ru.spcex.platform</groupId>
|
||||||
|
<artifactId>platform-messaging</artifactId>
|
||||||
|
</dependency>
|
||||||
|
<dependency>
|
||||||
|
<groupId>ru.spcex.platform</groupId>
|
||||||
|
<artifactId>platform-enum</artifactId>
|
||||||
|
</dependency>
|
||||||
|
<dependency>
|
||||||
|
<groupId>ru.spcex.platform</groupId>
|
||||||
|
<artifactId>platform-imdg-api-hazelcast-impl</artifactId>
|
||||||
|
</dependency>
|
||||||
|
<dependency>
|
||||||
|
<groupId>ru.spcex.clearing</groupId>
|
||||||
|
<artifactId>classes</artifactId>
|
||||||
|
</dependency>
|
||||||
|
<dependency>
|
||||||
|
<groupId>org.springframework.boot</groupId>
|
||||||
|
<artifactId>spring-boot-starter</artifactId>
|
||||||
|
</dependency>
|
||||||
|
<dependency>
|
||||||
|
<groupId>com.fasterxml.jackson.core</groupId>
|
||||||
|
<artifactId>jackson-databind</artifactId>
|
||||||
|
</dependency>
|
||||||
|
|
||||||
|
<!-- TEST -->
|
||||||
|
<dependency>
|
||||||
|
<groupId>org.springframework.boot</groupId>
|
||||||
|
<artifactId>spring-boot-starter-test</artifactId>
|
||||||
|
<scope>test</scope>
|
||||||
|
</dependency>
|
||||||
|
</dependencies>
|
||||||
|
|
||||||
|
<build>
|
||||||
|
<resources>
|
||||||
|
<resource>
|
||||||
|
<directory>src/main/resources</directory>
|
||||||
|
<excludes>
|
||||||
|
<exclude>application.properties</exclude>
|
||||||
|
</excludes>
|
||||||
|
<filtering>false</filtering>
|
||||||
|
</resource>
|
||||||
|
</resources>
|
||||||
|
<plugins>
|
||||||
|
<plugin>
|
||||||
|
<groupId>org.springframework.boot</groupId>
|
||||||
|
<artifactId>spring-boot-maven-plugin</artifactId>
|
||||||
|
<executions>
|
||||||
|
<execution>
|
||||||
|
<goals>
|
||||||
|
<goal>repackage</goal>
|
||||||
|
</goals>
|
||||||
|
</execution>
|
||||||
|
</executions>
|
||||||
|
<configuration>
|
||||||
|
<finalName>${project.artifactId}</finalName>
|
||||||
|
</configuration>
|
||||||
|
</plugin>
|
||||||
|
</plugins>
|
||||||
|
|
||||||
|
</build>
|
||||||
|
|
||||||
|
</project>
|
||||||
|
|
@ -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);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -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<String, Object> createConsumer(SchedulerServiceSettings settings) {
|
||||||
|
return KafkaConsumerFactory.consumer(settings.getKafkaConsumer());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
@Bean
|
||||||
|
public Producer<String, Object> createProducer(SchedulerServiceSettings settings) {
|
||||||
|
return KafkaProducerFactory.producer(settings.getKafkaProducer());
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
|
|
@ -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());
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
|
|
@ -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;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,33 @@
|
||||||
|
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.SchedulerAllToday;
|
||||||
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
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<SchedulerAllToday> schedulerAllTodayMap;
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
public SchedulerAllTodayService(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaProducer,
|
||||||
|
ImdgProvider imdgProvider) {
|
||||||
|
super(kafkaQueue, kafkaProducer);
|
||||||
|
this.schedulerAllTodayMap = imdgProvider.getImdg(IMDGDistributedNames.Map_SchedulerAllToday, SchedulerAllToday.class);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void afterPropertiesSet() {
|
||||||
|
init();
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
|
|
@ -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<KeyRate> keyRateMap;
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
public SchedulerService(Consumer<String, Object> kafkaQueue, Producer<String, Object> 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<KeyRateNewRequest> 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<KeyRateUpdateRequest> 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<CommonDeleteRequest> userRequest) {
|
||||||
|
CommonDeleteRequest req = userRequest.getRequestPayload();
|
||||||
|
log.debug("CommonDeleteRequest received id = {}", req.getId());
|
||||||
|
KeyRate keyRate = keyRateMap.getSingleObjectByID(req.getId());
|
||||||
|
keyRateMap.delete(keyRate);
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
|
|
@ -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
|
||||||
|
|
@ -0,0 +1,38 @@
|
||||||
|
<?xml version="1.0" encoding="UTF-8"?>
|
||||||
|
<configuration>
|
||||||
|
<appender name="CONSOLE" class="ch.qos.logback.core.ConsoleAppender">
|
||||||
|
<encoder>
|
||||||
|
<!-- |%X{ru.nbch.scoring.web.logging.mdc_key}-->
|
||||||
|
<Pattern>%date{HH:mm:ss.SSS} [%thread] %-5level %class{0}:%line - %message%n</Pattern>
|
||||||
|
<charset>utf-8</charset>
|
||||||
|
</encoder>
|
||||||
|
</appender>
|
||||||
|
<appender name="FILE" class="ch.qos.logback.core.rolling.RollingFileAppender">
|
||||||
|
<file>./logs/utility-service.log</file>
|
||||||
|
<encoder>
|
||||||
|
<!-- |%X{ru.nbch.scoring.web.logging.mdc_key}-->
|
||||||
|
<Pattern>%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %class{0}:%msg%n</Pattern>
|
||||||
|
<charset>utf8</charset>
|
||||||
|
</encoder>
|
||||||
|
<rollingPolicy class="ch.qos.logback.core.rolling.FixedWindowRollingPolicy">
|
||||||
|
<fileNamePattern>
|
||||||
|
./logs/utility-service.%i.log
|
||||||
|
</fileNamePattern>
|
||||||
|
<minIndex>1</minIndex>
|
||||||
|
<maxIndex>10</maxIndex>
|
||||||
|
</rollingPolicy>
|
||||||
|
<triggeringPolicy class="ch.qos.logback.core.rolling.SizeBasedTriggeringPolicy">
|
||||||
|
<maxFileSize>500MB</maxFileSize>
|
||||||
|
</triggeringPolicy>
|
||||||
|
</appender>
|
||||||
|
|
||||||
|
<root level="warn">
|
||||||
|
<appender-ref ref="CONSOLE"/>
|
||||||
|
<appender-ref ref="FILE"/>
|
||||||
|
</root>
|
||||||
|
|
||||||
|
<logger name="ru.spcex" level="debug" additivity="false">
|
||||||
|
<appender-ref ref="FILE"/>
|
||||||
|
<appender-ref ref="CONSOLE"/>
|
||||||
|
</logger>
|
||||||
|
</configuration>
|
||||||
|
|
@ -9,6 +9,10 @@ public interface Consts {
|
||||||
String DESTINATION_KEY_RATE_UPDATE = "key-rate-update";
|
String DESTINATION_KEY_RATE_UPDATE = "key-rate-update";
|
||||||
String DESTINATION_KEY_RATE_DELETE = "key-rate-delete";
|
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_DELETE = "company-delete";
|
||||||
String DESTINATION_COMPANY_INFO_UPDATE = "company-info-update";
|
String DESTINATION_COMPANY_INFO_UPDATE = "company-info-update";
|
||||||
String DESTINATION_COMPANY_SYMBOL_UPDATE = "company-symbol-update";
|
String DESTINATION_COMPANY_SYMBOL_UPDATE = "company-symbol-update";
|
||||||
|
|
|
||||||
|
|
@ -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;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -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;
|
||||||
|
}
|
||||||
|
}
|
||||||
Loading…
Add table
Reference in a new issue