From fa2cb814c806a4d189223cdae61d5c4d3f887ba4 Mon Sep 17 00:00:00 2001 From: AKurakin Date: Thu, 29 Jun 2023 11:55:18 +0300 Subject: [PATCH] =?UTF-8?q?utility-service=20=D0=BF=D0=BE=D0=BF=D1=80?= =?UTF-8?q?=D0=B0=D0=B2=D0=B8=D0=BB=20=D0=B2=D0=BE=D0=B7=D0=B2=D1=80=D0=B0?= =?UTF-8?q?=D1=89=D0=B5=D0=BD=D0=B8=D0=B5=20=D0=BE=D1=82=D0=B2=D0=B5=D1=82?= =?UTF-8?q?=D0=BE=D0=B2=20=D0=BF=D1=80=D0=B8=20=D0=BE=D0=B1=D1=80=D0=B0?= =?UTF-8?q?=D0=B1=D0=BE=D1=82=D0=BA=D0=B8=20Notification.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../utility/service/NotificationService.java | 25 +++++++++++-------- .../NotificationFeedbackRequest.java | 10 ++++++++ .../messaging/service/QueueConsumer.java | 6 +++++ 3 files changed, 30 insertions(+), 11 deletions(-) diff --git a/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/NotificationService.java b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/NotificationService.java index a8cec87fb..d47b0b6d6 100644 --- a/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/NotificationService.java +++ b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/NotificationService.java @@ -1,6 +1,7 @@ package ru.spcex.clearing.utility.service; import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; @@ -31,10 +32,10 @@ public class NotificationService extends QueueConsumer implements InitializingBe private final KafkaSender kafkaSender; @Autowired - public NotificationService(Consumer kafkaQueue, + public NotificationService(Consumer kafkaQueue, Producer kafkaProducer, ImdgProvider imdgProvider, KafkaSender kafkaSender) { - super(kafkaQueue); + super(kafkaQueue, kafkaProducer); this.notificationMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Notification, Notification.class); this.kafkaSender = kafkaSender; } @@ -51,17 +52,17 @@ public class NotificationService extends QueueConsumer implements InitializingBe } private void notificationNewRequest(BaseRequest notificationNewRequestBaseRequest) { - log.info("Starting notificationNewRequest processing..."); + log.info("Starting notificationNewRequest {} processing...", notificationNewRequestBaseRequest.getId()); NotificationNewRequest request = notificationNewRequestBaseRequest.getRequestPayload(); Notification notification = notificationBuilderFromNotificationNewRequest(request); notificationMap.insert(notification); - log.info("NotificationNewRequest successfully processed! Generated id:\t{}", notification.getId()); + log.info("NotificationNewRequest {} successfully processed! Generated id:\t{}", notificationNewRequestBaseRequest.getId(), notification.getId()); } private void notificationUpdateRequest(BaseRequest notificationUpdateRequestBaseRequest) { NotificationUpdateRequest request = notificationUpdateRequestBaseRequest.getRequestPayload(); - Long id = request.getId(); - log.info("Starting notificationUpdateRequest processing by id: {} ...", id); + Long id = request == null ? null : request.getId(); + log.info("Starting notificationUpdateRequest {} processing by id: {} ...", notificationUpdateRequestBaseRequest.getId(), id); Notification notification = notificationMap.getSingleObjectByID(id); if (notification == null) { throw new NullPointerException("Notification by id:{" + id + "} is null!"); @@ -69,8 +70,8 @@ public class NotificationService extends QueueConsumer implements InitializingBe notification.setNotificationStatus(request.getNotificationStatus()); notification.setUpdated(Instant.now()); notificationMap.update(notification); - sendFeedbackString(notification.getObjectType(), notification.getNotificationStatus()); - log.info("NotificationUpdateRequest successfully processed!"); + sendFeedbackString(notification); + log.info("NotificationUpdateRequest {} successfully processed!", notificationUpdateRequestBaseRequest.getId()); } private Notification notificationBuilderFromNotificationNewRequest(NotificationNewRequest request) { @@ -89,16 +90,18 @@ public class NotificationService extends QueueConsumer implements InitializingBe return notification; } - private void sendFeedbackString(String objectType, String status) { + private void sendFeedbackString(Notification notification) { + String objectType = notification.getObjectType(); if (objectType.equalsIgnoreCase(statement.getKey())) { - kafkaSender.sendRequestToQueue(CLEARING_NOTIFICATION_FEEDBACK, buildFeedbackRequest(status)); + kafkaSender.sendRequestToQueue(CLEARING_NOTIFICATION_FEEDBACK, buildFeedbackRequest(notification.getId(), notification.getNotificationStatus())); } else { log.warn("Unsupported notification type {}", objectType); } } - private NotificationFeedbackRequest buildFeedbackRequest(String status) { + private NotificationFeedbackRequest buildFeedbackRequest(Long notificationId, String status) { NotificationFeedbackRequest request = new NotificationFeedbackRequest(); + request.setNotificationId(notificationId); request.setNotificationStatus(status); return request; } diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/utilities/NotificationFeedbackRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/utilities/NotificationFeedbackRequest.java index f58461bfc..1af401322 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/utilities/NotificationFeedbackRequest.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/utilities/NotificationFeedbackRequest.java @@ -4,9 +4,19 @@ import com.fasterxml.jackson.annotation.JsonProperty; import com.fasterxml.jackson.databind.annotation.JsonSerialize; public class NotificationFeedbackRequest { + @JsonProperty + public Long notificationId; @JsonProperty public String notificationStatus; + public Long getNotificationId() { + return notificationId; + } + + public void setNotificationId(Long notificationId) { + this.notificationId = notificationId; + } + public String getNotificationStatus() { return notificationStatus; } 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 0e6d0ac3e..bfe4c6cb4 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 @@ -56,6 +56,9 @@ public class QueueConsumer implements AutoCloseable { this.outputExecutor = Executors.newSingleThreadExecutor(); this.json = new ObjectMapper(); this.supportStartOffsetTimeWindow = false; + if (producer == null) { + log.debug("For {} kafka producer not set. Did not send reply.", getClass().getName()); + } } /** @@ -68,6 +71,9 @@ public class QueueConsumer implements AutoCloseable { public QueueConsumer(Consumer kafkaQueue, Producer kafkaResponseQueue) { this(kafkaQueue); this.producer = kafkaResponseQueue; + if (producer == null) { + log.warn("For {} kafka producer not set. Did not send reply.", getClass().getName()); + } } protected boolean needsProcessing(String topicName, BaseRequest request) {