From 246dc403af9e68581c4ae9b7df948ee8d9ceefa7 Mon Sep 17 00:00:00 2001 From: ialbert Date: Wed, 3 May 2023 17:31:27 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-277 --- .../sender/KafkaSyncRequestReplySender.java | 78 ------------------- 1 file changed, 78 deletions(-) delete mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSyncRequestReplySender.java diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSyncRequestReplySender.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSyncRequestReplySender.java deleted file mode 100644 index 2d176d8e6..000000000 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/sender/KafkaSyncRequestReplySender.java +++ /dev/null @@ -1,78 +0,0 @@ -package ru.spcex.clearing.platform.messaging.service.sender; - -import org.apache.kafka.clients.producer.Producer; -import org.apache.kafka.clients.producer.ProducerRecord; -import org.apache.kafka.clients.producer.RecordMetadata; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; -import org.springframework.kafka.requestreply.ReplyingKafkaTemplate; -import ru.spcex.clearing.platform.messaging.domain.ActionType; -import ru.spcex.clearing.platform.messaging.domain.BaseRequest; -import ru.spcex.clearing.platform.messaging.service.RequestInfo; -import ru.spcex.platform.classes.base.SpcexObjectBase; -import ru.spcex.platform.utils.log.ExceptionUtils; - -import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.Future; -import java.util.function.Function; -import java.util.function.Supplier; - -public class KafkaSyncRequestReplySender { - private final Logger log = LoggerFactory.getLogger(getClass()); - private final Map allImdgMaps; - private Producer kafka; - private Supplier idGenerator; - private Function imdgGenerator; - private ReplyingKafkaTemplate kafkaTemplate; - - KafkaSyncRequestReplySender() { - this.allImdgMaps = new ConcurrentHashMap<>(); - } - - public static KafkaSenderBuilderImpl setup() { - return new KafkaSenderBuilderImpl(); - } - - - public Long sendRequestToQueue(String destination, Object requestPayload) { - BaseRequest request = new BaseRequest<>(); - request.setId(idGenerator.get()); - request.setActionType(ActionType.SYSTEM); - request.setRequestPayload(requestPayload); - //сохраняет данные о запросе в хранилище - saveRequestToStorage(destination, request); - Future send = kafka.send(new ProducerRecord<>(destination, request)); - try { - send.get(); - } catch (InterruptedException | ExecutionException e) { - log.error(ExceptionUtils.getStackTrace(e)); - return null; - } - return request.getId(); - } - - private void saveRequestToStorage(String destination, BaseRequest request) { - KafkaImdgInsert imdgInsert = getImdg(destination); - RequestInfo requestInfo = RequestInfo.create(request.getId()); - imdgInsert.insert(requestInfo); - } - - @SuppressWarnings("unchecked") - private KafkaImdgInsert getImdg(String mapName) { - return allImdgMaps.computeIfAbsent(mapName, (mapName1) -> imdgGenerator.apply(mapName)); - } - - void setProducer(Producer kafka) { - this.kafka = kafka; - } - - void setIdGenerator(Supplier idGenerator) { - this.idGenerator = idGenerator; - } - - void setImdgProvider(Function imdgGenerator) { - this.imdgGenerator = imdgGenerator; - } -}