diff --git a/clearing-parent/backend-api/pom.xml b/clearing-parent/backend-api/pom.xml
index a50260658..558019a31 100644
--- a/clearing-parent/backend-api/pom.xml
+++ b/clearing-parent/backend-api/pom.xml
@@ -38,8 +38,8 @@
platform-imdg-api-hazelcast-impl
- org.apache.kafka
- kafka-clients
+ ru.spcex.platform
+ platform-messaging
diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/KafkaConfig.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/KafkaConfig.java
index ee55cf976..3febf0247 100644
--- a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/KafkaConfig.java
+++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/KafkaConfig.java
@@ -1,14 +1,11 @@
package ru.spcex.clearing.backendapi.config;
-import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import ru.spcex.clearing.backendapi.config.element.BackendApiSettings;
-import ru.spcex.clearing.backendapi.config.element.KafkaSettings;
-
-import java.util.Properties;
+import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory;
@Configuration
public class KafkaConfig {
@@ -16,29 +13,7 @@ public class KafkaConfig {
@Autowired
@Bean
public Producer createProducer(BackendApiSettings settings) {
- KafkaSettings kafkaSettings = settings.getKafka();
- Properties kafkaProps = new Properties();
-
- //Assign localhost id
- kafkaProps.put("bootstrap.servers", kafkaSettings.getBootstrapServers());
-
- //Set acknowledgements for producer requests.
- kafkaProps.put("acks", kafkaSettings.getAcks());
-
- //If the request fails, the producer can automatically retry,
- kafkaProps.put("retries", kafkaSettings.getRetries());
-
- //Specify buffer size in config
- kafkaProps.put("batch.size", kafkaSettings.getBatchSize());
-
- //Reduce the no of requests less than 0
- kafkaProps.put("linger.ms", kafkaSettings.getLingerMs());
- //The buffer.memory controls the total amount of memory available to the producer for buffering.
- kafkaProps.put("buffer.memory", kafkaSettings.getBufferMemory());
- kafkaProps.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
- kafkaProps.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
-
- return new KafkaProducer<>(kafkaProps);
+ return KafkaProducerFactory.producer(settings.getKafka());
}
}
diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/element/BackendApiSettings.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/element/BackendApiSettings.java
index 4f521ef5f..f4180e7c2 100644
--- a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/element/BackendApiSettings.java
+++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/element/BackendApiSettings.java
@@ -3,6 +3,7 @@ package ru.spcex.clearing.backendapi.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.KafkaSettings;
import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams;
@Component
diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/cud/CudController.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/cud/CudController.java
new file mode 100644
index 000000000..7d7891d1c
--- /dev/null
+++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/cud/CudController.java
@@ -0,0 +1,86 @@
+package ru.spcex.clearing.backendapi.controller.cud;
+
+import com.fasterxml.jackson.annotation.JsonProperty;
+import com.fasterxml.jackson.core.JsonProcessingException;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import org.apache.kafka.clients.producer.Producer;
+import org.apache.kafka.clients.producer.ProducerRecord;
+import org.apache.kafka.clients.producer.RecordMetadata;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Controller;
+import org.springframework.web.bind.annotation.PathVariable;
+import org.springframework.web.bind.annotation.RequestBody;
+import org.springframework.web.bind.annotation.RequestMapping;
+import org.springframework.web.bind.annotation.RequestMethod;
+import ru.spcex.clearing.backendapi.controller.response.BasicSpcexResponse;
+import ru.spcex.clearing.backendapi.domain.actions.IAction;
+import ru.spcex.clearing.backendapi.meta.CudMetaService;
+
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.Future;
+
+@Controller
+@RequestMapping("/cud")
+public class CudController {
+ private final Producer kafka;
+ private final ObjectMapper json;
+ private final CudMetaService meta;
+
+ @Autowired
+ public CudController(Producer kafka, CudMetaService meta) {
+ this.kafka = kafka;
+ this.meta = meta;
+ this.json = new ObjectMapper();
+ }
+
+ @RequestMapping(value = "/{destination}", method = RequestMethod.POST)
+ public BasicSpcexResponse add(@PathVariable("destination") String destination,
+ @RequestBody String body) throws JsonProcessingException, ExecutionException, InterruptedException {
+ Class> actionClazz = meta.byDestination(destination);
+ if (actionClazz == null) {
+ throw new UnsupportedOperationException("unsupported destination " + destination);
+ }
+ IAction> iAction = json.readValue(body, actionClazz);
+
+ Future send = kafka.send(new ProducerRecord<>(destination, body));
+ RecordMetadata kafkaMetaData = send.get();
+ MetaDataResponse responseToClient = new MetaDataResponse();
+ responseToClient.setOffset(kafkaMetaData.offset());
+ responseToClient.setPartition(kafkaMetaData.partition());
+ responseToClient.setTopic(kafkaMetaData.topic());
+ return responseToClient;
+ }
+
+ private static class MetaDataResponse extends BasicSpcexResponse {
+ @JsonProperty
+ private Long offset;
+ @JsonProperty
+ private Integer partition;
+ @JsonProperty
+ private String topic;
+
+ public Long getOffset() {
+ return offset;
+ }
+
+ public void setOffset(Long offset) {
+ this.offset = offset;
+ }
+
+ public Integer getPartition() {
+ return partition;
+ }
+
+ public void setPartition(Integer partition) {
+ this.partition = partition;
+ }
+
+ public String getTopic() {
+ return topic;
+ }
+
+ public void setTopic(String topic) {
+ this.topic = topic;
+ }
+ }
+}
diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/request/cud/MoneyMarketCreate.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/request/cud/MoneyMarketCreate.java
new file mode 100644
index 000000000..ac90c6039
--- /dev/null
+++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/request/cud/MoneyMarketCreate.java
@@ -0,0 +1,96 @@
+package ru.spcex.clearing.backendapi.controller.request.cud;
+
+import com.fasterxml.jackson.annotation.JsonFormat;
+import com.fasterxml.jackson.annotation.JsonProperty;
+import ru.spcex.clearing.backendapi.domain.actions.IAction;
+
+import java.time.Instant;
+
+public class MoneyMarketCreate implements IAction