parent
6bb00569d3
commit
333d935a26
2 changed files with 20 additions and 6 deletions
|
|
@ -13,6 +13,7 @@ import org.springframework.kafka.core.ProducerFactory;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory;
|
import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory;
|
||||||
import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory;
|
import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory;
|
||||||
|
import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings;
|
||||||
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
|
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
|
||||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||||
import ru.spcex.clearing.utility.config.settings.UtilityServiceSettings;
|
import ru.spcex.clearing.utility.config.settings.UtilityServiceSettings;
|
||||||
|
|
@ -35,6 +36,11 @@ public class KafkaConfig {
|
||||||
return KafkaProducerFactory.producer(settings.getKafkaProducer());
|
return KafkaProducerFactory.producer(settings.getKafkaProducer());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Bean
|
||||||
|
public ProducerFactory<String, Object> pf(UtilityServiceSettings settings) {
|
||||||
|
KafkaProducerSettings kafkaSettings = settings.getKafkaProducer();
|
||||||
|
return KafkaProducerFactory.producerFactory(kafkaSettings);
|
||||||
|
}
|
||||||
|
|
||||||
@Bean("kafkaTemplate")
|
@Bean("kafkaTemplate")
|
||||||
public KafkaTemplate<String, Object> kafkaTemplate(ProducerFactory<String, Object> pf) {
|
public KafkaTemplate<String, Object> kafkaTemplate(ProducerFactory<String, Object> pf) {
|
||||||
|
|
|
||||||
|
|
@ -9,7 +9,9 @@ import org.springframework.stereotype.Service;
|
||||||
import ru.clearing.classes.statics.data.misc.Notification;
|
import ru.clearing.classes.statics.data.misc.Notification;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.*;
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationFeedbackRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationNewRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationUpdateRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||||
import ru.spcex.platform.imdg.api.Imdg;
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
|
|
@ -57,9 +59,13 @@ public class NotificationService extends QueueConsumer implements InitializingBe
|
||||||
}
|
}
|
||||||
|
|
||||||
private void notificationUpdateRequest(BaseRequest<NotificationUpdateRequest> notificationUpdateRequestBaseRequest) {
|
private void notificationUpdateRequest(BaseRequest<NotificationUpdateRequest> notificationUpdateRequestBaseRequest) {
|
||||||
log.info("Starting notificationUpdateRequest processing...");
|
|
||||||
NotificationUpdateRequest request = notificationUpdateRequestBaseRequest.getRequestPayload();
|
NotificationUpdateRequest request = notificationUpdateRequestBaseRequest.getRequestPayload();
|
||||||
Notification notification = notificationMap.getSingleObjectByID(request.getId());
|
Long id = request.getId();
|
||||||
|
log.info("Starting notificationUpdateRequest processing by id: {} ...", id);
|
||||||
|
Notification notification = notificationMap.getSingleObjectByID(id);
|
||||||
|
if (notification == null) {
|
||||||
|
throw new NullPointerException("Notification by id:{" + id + "} is null!");
|
||||||
|
}
|
||||||
notification.setNotificationStatus(request.getNotificationStatus());
|
notification.setNotificationStatus(request.getNotificationStatus());
|
||||||
notification.setUpdated(Instant.now());
|
notification.setUpdated(Instant.now());
|
||||||
notificationMap.update(notification);
|
notificationMap.update(notification);
|
||||||
|
|
@ -80,13 +86,15 @@ public class NotificationService extends QueueConsumer implements InitializingBe
|
||||||
return notification;
|
return notification;
|
||||||
}
|
}
|
||||||
|
|
||||||
private void sendFeedbackString(String objectType, String status){
|
private void sendFeedbackString(String objectType, String status) {
|
||||||
if (objectType.equalsIgnoreCase(statement.getKey())){
|
if (objectType.equalsIgnoreCase(statement.getKey())) {
|
||||||
kafkaSender.sendRequestToQueue(CLEARING_NOTIFICATION_FEEDBACK, buildFeedbackRequest(status));
|
kafkaSender.sendRequestToQueue(CLEARING_NOTIFICATION_FEEDBACK, buildFeedbackRequest(status));
|
||||||
|
} else {
|
||||||
|
log.warn("Unsupported notification type {}", objectType);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private NotificationFeedbackRequest buildFeedbackRequest(String status){
|
private NotificationFeedbackRequest buildFeedbackRequest(String status) {
|
||||||
NotificationFeedbackRequest request = new NotificationFeedbackRequest();
|
NotificationFeedbackRequest request = new NotificationFeedbackRequest();
|
||||||
request.setNotificationStatus(status);
|
request.setNotificationStatus(status);
|
||||||
return request;
|
return request;
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue