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;
- }
- }
}