diff --git a/clearing-parent/balance-service/pom.xml b/clearing-parent/balance-service/pom.xml new file mode 100644 index 000000000..60e35a567 --- /dev/null +++ b/clearing-parent/balance-service/pom.xml @@ -0,0 +1,68 @@ + + + + clearing-parent + ru.spcex.clearing + SPCEX-1.0.0.0 + + 4.0.0 + + balance-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/balance-service/src/main/java/ru/spcex/clearing/balance/BalanceServiceApplication.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/BalanceServiceApplication.java new file mode 100644 index 000000000..d4fb58bd9 --- /dev/null +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/BalanceServiceApplication.java @@ -0,0 +1,12 @@ +package ru.spcex.clearing.balance; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +@SpringBootApplication +public class BalanceServiceApplication { + public static void main(String[] args) { + SpringApplication app = new SpringApplication(BalanceServiceApplication.class); + app.run(args); + } +} diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/BalanceImdgConfig.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/BalanceImdgConfig.java new file mode 100644 index 000000000..a4ffd6dd2 --- /dev/null +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/BalanceImdgConfig.java @@ -0,0 +1,47 @@ +package ru.spcex.clearing.balance.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.balance.config.element.BalanceServiceSettings; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; + +@Configuration +public class BalanceImdgConfig { + 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, + BalanceServiceSettings settings + ) { + return new HazelcastService(taskExecutorHazelcastClientInitializer, + taskExecutorIdGeneratorAwaiter, + settings.getHazelcast()); + } + +} diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/KafkaConfig.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/KafkaConfig.java new file mode 100644 index 000000000..1cfeb4a56 --- /dev/null +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/KafkaConfig.java @@ -0,0 +1,20 @@ +package ru.spcex.clearing.balance.config; + +import org.apache.kafka.clients.consumer.Consumer; +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.balance.config.element.BalanceServiceSettings; +import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; + +@Configuration +public class KafkaConfig { + @Autowired + @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) + @Bean + public Consumer createProducer(BalanceServiceSettings settings) { + return KafkaConsumerFactory.consumer(settings.getKafka()); + } +} diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/element/BalanceServiceSettings.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/element/BalanceServiceSettings.java new file mode 100644 index 000000000..dbdb8c45a --- /dev/null +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/element/BalanceServiceSettings.java @@ -0,0 +1,31 @@ +package ru.spcex.clearing.balance.config.element; + +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("balance-service") +public class BalanceServiceSettings { + 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/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf02Service.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf02Service.java new file mode 100644 index 000000000..3a9b1b2e5 --- /dev/null +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf02Service.java @@ -0,0 +1,59 @@ +package ru.spcex.clearing.balance.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.clearing.classes.statics.data.sdf.SDf02; +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.balance.SDf02NewRequest; +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 Sdf02Service extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg sdf8Map; + + + public Sdf02Service(Consumer kafkaQueue, ImdgProvider imdgProvider) { + super(kafkaQueue); + this.sdf8Map = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf02, SDf02.class); + } + + @Override + public void afterPropertiesSet() { + callback(SDf02NewRequest.class) + .setConsumer(this::newSDf02) + .forDestination(Consts.DESTINATION_SDF02_NEW, callbacks::put); + init(); + } + + private void newSDf02(BaseRequest userRequest) { + SDf02NewRequest req = userRequest.getRequestPayload(); + log.debug("SDf02NewRequest received"); + SDf02 sDf02 = new SDf02(); + sDf02.setCurr_code(req.getCurr_code()); + sDf02.setAccount(req.getAccount()); + sDf02.setRemainder(req.getRemainder()); + sDf02.setDeal(req.getDeal()); + sDf02.setAcc_code(req.getAcc_code()); + sDf02.setDat(req.getDat()); + sDf02.setMarket(req.getMarket()); + sDf02.setAcc_name(req.getAcc_name()); + sDf02.setAcc_type(req.getAcc_type()); + sDf02.setSumengage(req.getSumengage()); + sDf02.setSumunblock(req.getSumunblock()); + sDf02.setFile_type(req.getFile_type()); + sDf02.setResult(req.getResult()); + sDf02.setFileName(req.getFileName()); + sDf02.setGenerationTime(req.getGenerationTime()); + sDf02.setGenerationId(req.getGenerationId()); + sdf8Map.insert(sDf02); + log.debug("successfully processed, new id {}", sDf02.getId()); + } +} diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf08Service.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf08Service.java new file mode 100644 index 000000000..52e56b56a --- /dev/null +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf08Service.java @@ -0,0 +1,48 @@ +package ru.spcex.clearing.balance.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.clearing.classes.statics.data.sdf.SDf08; +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.balance.SDf08NewRequest; +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 Sdf08Service extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg sdf8Map; + + + public Sdf08Service(Consumer kafkaQueue, ImdgProvider imdgProvider) { + super(kafkaQueue); + this.sdf8Map = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf08, SDf08.class); + } + + @Override + public void afterPropertiesSet() { + callback(SDf08NewRequest.class) + .setConsumer(this::newSDf08) + .forDestination(Consts.DESTINATION_SDF08_NEW, callbacks::put); + init(); + } + + private void newSDf08(BaseRequest userRequest) { + SDf08NewRequest req = userRequest.getRequestPayload(); + log.debug("SDf08NewRequest received"); + SDf08 sDf08 = new SDf08(); + sDf08.setNumber(req.getNumber()); + sDf08.setDatetime(req.getDatetime()); + sDf08.setFileName(req.getFileName()); + sDf08.setGenerationTime(req.getGenerationTime()); + sDf08.setGenerationId(req.getGenerationId()); + sdf8Map.insert(sDf08); + log.debug("successfully processed, new id {}", sDf08.getId()); + } +} diff --git a/clearing-parent/balance-service/src/main/resources/application.properties b/clearing-parent/balance-service/src/main/resources/application.properties new file mode 100644 index 000000000..2e1bbbab3 --- /dev/null +++ b/clearing-parent/balance-service/src/main/resources/application.properties @@ -0,0 +1,12 @@ +spring.main.web-application-type=none +balance-service.hazelcast.cluster-members=127.0.0.1 +balance-service.hazelcast.login=dev +balance-service.hazelcast.password=dev-pass + +balance-service.kafka.bootstrap-servers=localhost:9092 +balance-service.kafka.group-id=dev-group +balance-service.kafka.enable-auto-commit=false +balance-service.kafka.session-timeout-ms=30000 +balance-service.kafka.auto-offset-reset=latest +balance-service.kafka.linger-ms=1 +balance-service.kafka.buffer-memory=33554432 \ No newline at end of file diff --git a/clearing-parent/balance-service/src/main/resources/logback.xml b/clearing-parent/balance-service/src/main/resources/logback.xml new file mode 100644 index 000000000..7b16a2d56 --- /dev/null +++ b/clearing-parent/balance-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/balance-service.log + + + %d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %class{0}:%msg%n + utf8 + + + + ./logs/balance-service.%i.log + + 1 + 10 + + + 500MB + + + + + + + + + + + + + diff --git a/clearing-parent/pom.xml b/clearing-parent/pom.xml index 859773303..2ef81dddd 100644 --- a/clearing-parent/pom.xml +++ b/clearing-parent/pom.xml @@ -28,6 +28,7 @@ company-service reports-service account-service + balance-service 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 52ddf032d..98532fa97 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 @@ -14,8 +14,12 @@ public interface Consts { String DESTINATION_COMPANY_SYMBOL_UPDATE = "company-symbol-update"; String DESTINATION_CONTACT_UPDATE = "contact-update"; String DESTINATION_CLEARING_MEMBER_CATEGORY_UPDATE = "clearing-member-category-update"; + String DESTINATION_CLEARING_MEMBER_CATEGORY_DELETE = "clearing-member-category-delete"; + String DESTINATION_BANK_ACCOUNT_DELETE = "bank-account-delete"; String DESTINATION_BANK_ACCOUNT_UPDATE = "bank-account-update"; String DESTINATION_BANK_ACCOUNT_NEW = "bank-account-new"; - String DESTINATION_CLEARING_MEMBER_CATEGORY_DELETE = "clearing-member-category-delete"; + + String DESTINATION_SDF08_NEW = "s-df-08-new"; + String DESTINATION_SDF02_NEW = "s-df-02-new"; } diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/balance/SDf02NewRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/balance/SDf02NewRequest.java new file mode 100644 index 000000000..1ca2a31c8 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/balance/SDf02NewRequest.java @@ -0,0 +1,174 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.balance; + +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.InstantDeserializer; +import ru.spcex.clearing.platform.messaging.domain.json.serialize.InstantSerializer; + +import java.time.Instant; + +public class SDf02NewRequest { + @JsonProperty + private String curr_code; + @JsonProperty + private String account; + @JsonProperty + private String remainder; + @JsonProperty + private String deal; + @JsonProperty + private String acc_code; + @JsonProperty + private String dat; + @JsonProperty + private String market; + @JsonProperty + private String acc_name; + @JsonProperty + private String acc_type; + @JsonProperty + private String sumengage; + @JsonProperty + private String sumunblock; + @JsonProperty + private String file_type; + @JsonProperty + private String result; + @JsonProperty + private String fileName; + @JsonSerialize(using = InstantSerializer.class) + @JsonDeserialize(using = InstantDeserializer.class) + @JsonProperty + private Instant generationTime; + @JsonProperty + private Long generationId; + + public String getCurr_code() { + return curr_code; + } + + public void setCurr_code(String curr_code) { + this.curr_code = curr_code; + } + + public String getAccount() { + return account; + } + + public void setAccount(String account) { + this.account = account; + } + + public String getRemainder() { + return remainder; + } + + public void setRemainder(String remainder) { + this.remainder = remainder; + } + + public String getDeal() { + return deal; + } + + public void setDeal(String deal) { + this.deal = deal; + } + + public String getAcc_code() { + return acc_code; + } + + public void setAcc_code(String acc_code) { + this.acc_code = acc_code; + } + + public String getDat() { + return dat; + } + + public void setDat(String dat) { + this.dat = dat; + } + + public String getMarket() { + return market; + } + + public void setMarket(String market) { + this.market = market; + } + + public String getAcc_name() { + return acc_name; + } + + public void setAcc_name(String acc_name) { + this.acc_name = acc_name; + } + + public String getAcc_type() { + return acc_type; + } + + public void setAcc_type(String acc_type) { + this.acc_type = acc_type; + } + + public String getSumengage() { + return sumengage; + } + + public void setSumengage(String sumengage) { + this.sumengage = sumengage; + } + + public String getSumunblock() { + return sumunblock; + } + + public void setSumunblock(String sumunblock) { + this.sumunblock = sumunblock; + } + + public String getFile_type() { + return file_type; + } + + public void setFile_type(String file_type) { + this.file_type = file_type; + } + + public String getResult() { + return result; + } + + public void setResult(String result) { + this.result = result; + } + + public String getFileName() { + return fileName; + } + + public void setFileName(String fileName) { + this.fileName = fileName; + } + + public Instant getGenerationTime() { + return generationTime; + } + + public void setGenerationTime(Instant generationTime) { + this.generationTime = generationTime; + } + + public Long getGenerationId() { + return generationId; + } + + public void setGenerationId(Long generationId) { + this.generationId = generationId; + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/balance/SDf08NewRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/balance/SDf08NewRequest.java new file mode 100644 index 000000000..b0038e7f9 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/balance/SDf08NewRequest.java @@ -0,0 +1,67 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.balance; + +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.InstantDeserializer; +import ru.spcex.clearing.platform.messaging.domain.json.serialize.InstantSerializer; + +import java.math.BigDecimal; +import java.time.Instant; + +public class SDf08NewRequest { + @JsonProperty + private BigDecimal number; + @JsonSerialize(using = InstantSerializer.class) + @JsonDeserialize(using = InstantDeserializer.class) + @JsonProperty + private Instant datetime; + @JsonProperty + private String fileName; + @JsonSerialize(using = InstantSerializer.class) + @JsonDeserialize(using = InstantDeserializer.class) + @JsonProperty + private Instant generationTime; + @JsonProperty + private Long generationId; + + public BigDecimal getNumber() { + return number; + } + + public void setNumber(BigDecimal number) { + this.number = number; + } + + public Instant getDatetime() { + return datetime; + } + + public void setDatetime(Instant datetime) { + this.datetime = datetime; + } + + public String getFileName() { + return fileName; + } + + public void setFileName(String fileName) { + this.fileName = fileName; + } + + public Instant getGenerationTime() { + return generationTime; + } + + public void setGenerationTime(Instant generationTime) { + this.generationTime = generationTime; + } + + public Long getGenerationId() { + return generationId; + } + + public void setGenerationId(Long generationId) { + this.generationId = generationId; + } +}