From a8c6ca094512ad5b6de016ad790b576e049e4a8f Mon Sep 17 00:00:00 2001 From: ialbert Date: Fri, 12 Aug 2022 13:17:55 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-16 --- clearing-parent/imdg-util/pom.xml | 20 -------- clearing-parent/imdg/pom.xml | 7 --- clearing-parent/pom.xml | 1 - clearing-parent/securities-service/pom.xml | 4 ++ .../cud/MoneyMarketSecurityService.java | 28 +++++++++-- .../iml/hazelcast/adapter/ImdgHazelcast.java | 16 ++++++- .../service/HazelcastServiceBase.java | 4 ++ .../clearing/imdg/IMDGDistributedNames.java | 0 .../messaging/logic/functional/Builder.java | 5 ++ .../logic/functional/BuilderConsumerStep.java | 7 +++ .../functional/BuilderDestinationStep.java | 7 +++ .../functional/ConsumerSpecificClass.java | 47 +++++++++++++++++++ .../messaging/service/QueueConsumer.java | 47 ++++--------------- 13 files changed, 122 insertions(+), 71 deletions(-) delete mode 100644 clearing-parent/imdg-util/pom.xml rename {clearing-parent/imdg-util => platform-parent/platform-imdg-api}/src/main/java/ru/spcex/clearing/imdg/IMDGDistributedNames.java (100%) create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/Builder.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/BuilderConsumerStep.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/BuilderDestinationStep.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/ConsumerSpecificClass.java diff --git a/clearing-parent/imdg-util/pom.xml b/clearing-parent/imdg-util/pom.xml deleted file mode 100644 index 184b5b734..000000000 --- a/clearing-parent/imdg-util/pom.xml +++ /dev/null @@ -1,20 +0,0 @@ - - - - clearing-parent - ru.spcex.clearing - SPCEX-1.0.0.0 - - 4.0.0 - - imdg-util - Клиентские утилиты для подключения к IMDG хранилищу - - - 17 - 17 - - - \ No newline at end of file diff --git a/clearing-parent/imdg/pom.xml b/clearing-parent/imdg/pom.xml index ea64f5b00..b3bd2cfea 100644 --- a/clearing-parent/imdg/pom.xml +++ b/clearing-parent/imdg/pom.xml @@ -42,13 +42,6 @@ compile - - ru.spcex.clearing - imdg-util - SPCEX-1.0.0.0 - compile - - diff --git a/clearing-parent/pom.xml b/clearing-parent/pom.xml index f1430a80d..3c0360f3e 100644 --- a/clearing-parent/pom.xml +++ b/clearing-parent/pom.xml @@ -20,7 +20,6 @@ classes backend-api imdg - imdg-util db-scripts dbf-importer dbf-exporter diff --git a/clearing-parent/securities-service/pom.xml b/clearing-parent/securities-service/pom.xml index 541d07d98..1c2aaf8b3 100644 --- a/clearing-parent/securities-service/pom.xml +++ b/clearing-parent/securities-service/pom.xml @@ -24,6 +24,10 @@ ru.spcex.platform platform-imdg-api-hazelcast-impl + + ru.spcex.clearing + classes + org.springframework.boot spring-boot-starter diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java index 45544ab74..4ca5c5e4c 100644 --- a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java @@ -4,24 +4,44 @@ import org.apache.kafka.clients.consumer.Consumer; import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; +import ru.clearing.classes.StaticData.Misc.MoneyMarketSecurity; +import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketCreateRequest; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.math.BigDecimal; @Service public class MoneyMarketSecurityService extends QueueConsumer implements InitializingBean { + private final Imdg moneyMarketSecurityMap; @Autowired - public MoneyMarketSecurityService(Consumer kafkaQueue) { + public MoneyMarketSecurityService(Consumer kafkaQueue, ImdgProvider imdgProvider) { super(kafkaQueue); + this.moneyMarketSecurityMap = imdgProvider.getImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, MoneyMarketSecurity.class); } @Override public void afterPropertiesSet() { - callback(Consts.DESTINATION_MONEY_MARKET_SECURITY_NEW, MoneyMarketCreateRequest.class) - .define(this::newMoneyMarket); + callback(MoneyMarketCreateRequest.class) + .setConsumer(this::newMoneyMarket) + .setDestination(Consts.DESTINATION_MONEY_MARKET_SECURITY_NEW, callbacks::put) + .build(); + init(); } private void newMoneyMarket(MoneyMarketCreateRequest req) { - + MoneyMarketSecurity mms = new MoneyMarketSecurity(); + mms.setStartDate(req.getStartDate()); + mms.setEndDate(req.getEndDate()); + mms.setNominalValue(BigDecimal.valueOf(req.getNominalValue())); + mms.setNominalCurrency(req.getNominalCurrencyId()); + mms.setInstrumentType(req.getInstrumentType()); + mms.setFullName(req.getFullName()); + mms.setSecuritySymbol(req.getSecuritySymbol()); + moneyMarketSecurityMap.insert(mms); +// (req.getLotSize()); } } diff --git a/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/adapter/ImdgHazelcast.java b/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/adapter/ImdgHazelcast.java index 34b201a88..d60cadcdd 100644 --- a/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/adapter/ImdgHazelcast.java +++ b/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/adapter/ImdgHazelcast.java @@ -1,10 +1,12 @@ package ru.spcex.platform.imdg.iml.hazelcast.adapter; import com.hazelcast.core.IMap; +import com.hazelcast.core.IdGenerator; import ru.spcex.platform.classes.base.SpcexObjectBase; import ru.spcex.platform.imdg.api.Imdg; public class ImdgHazelcast implements Imdg { + private IdGenerator idGenerator; private IMap map; @@ -15,5 +17,17 @@ public class ImdgHazelcast implements Imdg { public void setMap(IMap map) { this.map = map; } - //todo implement methods + + public IdGenerator getIdGenerator() { + return idGenerator; + } + + public void setIdGenerator(IdGenerator idGenerator) { + this.idGenerator = idGenerator; + } + + @Override + public void insert(T paramT) { + map.put(idGenerator.newId(), paramT); + } } diff --git a/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/service/HazelcastServiceBase.java b/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/service/HazelcastServiceBase.java index 93c026724..850ddbf3d 100644 --- a/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/service/HazelcastServiceBase.java +++ b/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/service/HazelcastServiceBase.java @@ -3,11 +3,13 @@ package ru.spcex.platform.imdg.iml.hazelcast.service; import com.hazelcast.client.HazelcastClient; import com.hazelcast.client.config.ClientConfig; import com.hazelcast.core.HazelcastInstance; +import com.hazelcast.core.IdGenerator; import com.hazelcast.core.LifecycleEvent; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.core.env.Environment; import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; +import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.platform.classes.base.SpcexObjectBase; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; @@ -214,7 +216,9 @@ public abstract class HazelcastServiceBase @Override public Imdg getImdg(String key, Class clazz) { ImdgHazelcast imdg = new ImdgHazelcast<>(); + IdGenerator generator = hazelcastInstance.getIdGenerator(IMDGDistributedNames.MAP_SEQUENCE_NAME); imdg.setMap(getHazelcast().getMap(key)); + imdg.setIdGenerator(generator); return imdg; } diff --git a/clearing-parent/imdg-util/src/main/java/ru/spcex/clearing/imdg/IMDGDistributedNames.java b/platform-parent/platform-imdg-api/src/main/java/ru/spcex/clearing/imdg/IMDGDistributedNames.java similarity index 100% rename from clearing-parent/imdg-util/src/main/java/ru/spcex/clearing/imdg/IMDGDistributedNames.java rename to platform-parent/platform-imdg-api/src/main/java/ru/spcex/clearing/imdg/IMDGDistributedNames.java diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/Builder.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/Builder.java new file mode 100644 index 000000000..e98db354f --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/Builder.java @@ -0,0 +1,5 @@ +package ru.spcex.clearing.platform.messaging.logic.functional; + +public interface Builder { + ConsumerSpecificClass build(); +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/BuilderConsumerStep.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/BuilderConsumerStep.java new file mode 100644 index 000000000..3fd881ecc --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/BuilderConsumerStep.java @@ -0,0 +1,7 @@ +package ru.spcex.clearing.platform.messaging.logic.functional; + +import java.util.function.Consumer; + +public interface BuilderConsumerStep { + BuilderDestinationStep setConsumer(Consumer consumer); +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/BuilderDestinationStep.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/BuilderDestinationStep.java new file mode 100644 index 000000000..d5b38d6af --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/BuilderDestinationStep.java @@ -0,0 +1,7 @@ +package ru.spcex.clearing.platform.messaging.logic.functional; + +import java.util.function.BiConsumer; + +public interface BuilderDestinationStep { + Builder setDestination(String destination, BiConsumer> mappingFun); +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/ConsumerSpecificClass.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/ConsumerSpecificClass.java new file mode 100644 index 000000000..3ffe10b30 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/ConsumerSpecificClass.java @@ -0,0 +1,47 @@ +package ru.spcex.clearing.platform.messaging.logic.functional; + +import java.util.function.BiConsumer; +import java.util.function.Consumer; + +public class ConsumerSpecificClass implements BuilderConsumerStep, BuilderDestinationStep, Builder { + private Class clazz; + private java.util.function.Consumer consumer; + + public ConsumerSpecificClass(Class clazz) { + this.clazz = clazz; + } + + public Builder setDestination(String destination, BiConsumer> fun) { + fun.accept(destination, this); + return this; + } + + public void accept(T obj) { + this.consumer.accept(obj); + } + + public void acceptRaw(Object obj) { + this.consumer.accept((T) obj); + } + + public Class getClazz() { + return clazz; + } + + public static BuilderConsumerStep build(Class clazz) { + return new ConsumerSpecificClass<>(clazz); + } + + @Override + public BuilderDestinationStep setConsumer(Consumer consumer) { + this.consumer = consumer; + return this; + } + + @Override + public ConsumerSpecificClass build() { + return this; + } + + +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java index 084b3dfce..f3186c7a2 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java @@ -6,6 +6,8 @@ import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.common.errors.WakeupException; +import ru.spcex.clearing.platform.messaging.logic.functional.BuilderConsumerStep; +import ru.spcex.clearing.platform.messaging.logic.functional.ConsumerSpecificClass; import java.time.Duration; import java.time.temporal.ChronoUnit; @@ -19,7 +21,7 @@ public class QueueConsumer implements AutoCloseable { private final AtomicBoolean closed = new AtomicBoolean(false); private final Consumer consumer; private final ExecutorService executor; - private final Map> callbacks; + protected final Map> callbacks; private final ObjectMapper json; public QueueConsumer(Consumer kafkaQueue) { @@ -29,18 +31,14 @@ public class QueueConsumer implements AutoCloseable { this.json = new ObjectMapper(); } - protected void addCallBack(String destination, ConsumerWithClass callback) { - this.callbacks.put(destination, callback); - } - - public void init() throws Exception { + public void init() { executor.submit(() -> { try { consumer.subscribe(callbacks.keySet()); while (!closed.get()) { ConsumerRecords records = consumer.poll(Duration.of(10, ChronoUnit.SECONDS)); for (ConsumerRecord next : records) { - ConsumerWithClass callback = callbacks.get(next.topic()); + ConsumerSpecificClass callback = callbacks.get(next.topic()); Object o = json.readValue((String) next.value(), callback.getClazz()); callback.acceptRaw(o); } @@ -55,42 +53,15 @@ public class QueueConsumer implements AutoCloseable { }); } + protected BuilderConsumerStep callback(Class clazz) { + return ConsumerSpecificClass.build(clazz); + } + @Override public void close() throws Exception { closed.set(true); consumer.wakeup(); } - protected ConsumerWithClass callback(String destination, Class clazz) { - ConsumerWithClass callback = new ConsumerWithClass<>(clazz); - this.callbacks.put(destination, callback); - return callback; - } - protected static class ConsumerWithClass { - private final Class clazz; - private java.util.function.Consumer consumer; - private ConsumerWithClass(Class clazz) { - this.clazz = clazz; - } - - - - public ConsumerWithClass define(java.util.function.Consumer consumer) { - this.consumer = consumer; - return this; - } - - public void accept(T obj) { - this.consumer.accept(obj); - } - - public void acceptRaw(Object obj) { - this.consumer.accept((T) obj); - } - - public Class getClazz() { - return clazz; - } - } }