From 723f177c25e15b75821eebcdb290d7e69eed98fc Mon Sep 17 00:00:00 2001 From: etreshenkov Date: Thu, 5 Oct 2023 19:31:54 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-562 --- .../utility/service/SdfPairChecker.java | 88 +++++++++++++++++++ .../ru/spcex/platform/enumeration/Task.java | 1 + .../spcex/platform/utils/time/TimeUtil.java | 13 +++ 3 files changed, 102 insertions(+) create mode 100644 clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/SdfPairChecker.java diff --git a/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/SdfPairChecker.java b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/SdfPairChecker.java new file mode 100644 index 000000000..4a256c0e8 --- /dev/null +++ b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/SdfPairChecker.java @@ -0,0 +1,88 @@ +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; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.sdf.SDf02; +import ru.clearing.classes.statics.data.sdf.SDf09; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationNewRequest; +import ru.spcex.clearing.platform.messaging.serialization.LogFormatter; +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.enumeration.Priority; +import ru.spcex.platform.enumeration.Task; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.api.predicate.ImdgPredicate; +import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.predicate.ImdgPredicateBuilderHazelcast; +import ru.spcex.platform.utils.time.TimeUtil; + +import java.time.Instant; +import java.time.temporal.ChronoUnit; +import java.util.Collection; + +@Service +public class SdfPairChecker extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg sDf02Imdg; + private final Imdg sDf09Imdg; + private final ImdgProvider imdgProvider; + private final KafkaSender kafkaSender; + + @Autowired + public SdfPairChecker(ImdgProvider imdgProvider, + KafkaSender kafkaSender, + Producer kafkaProducer, + Consumer kafkaQueue) { + super(kafkaQueue, kafkaProducer); + this.imdgProvider = imdgProvider; + this.kafkaSender = kafkaSender; + this.sDf02Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf02, SDf02.class); + this.sDf09Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf09, SDf09.class); + } + + @Override + public void afterPropertiesSet() throws Exception { + callback(LauncherCommandRequest.class) + .setConsumer(this::checkOnDate) + .forDestination(Task.checkSdfPair.topic(), callbacks::put); + init(); + } + + public void checkOnDate(BaseRequest authEvent) { + ImdgPredicateBuilder builder = ImdgPredicateBuilderHazelcast.instance(); + Instant now = Instant.now(); + Instant tradingDayStart = TimeUtil.startOfDay(now); + Instant tradingDayEnd = tradingDayStart.plus(1, ChronoUnit.DAYS); + ImdgPredicate predicate = builder.and( + builder.greatEqual("generationTime", tradingDayStart), + builder.less("generationTime", tradingDayEnd) + ); + Collection sdf09OnDate = sDf09Imdg.getCollectionObjectsByPredicate(predicate); + Collection sdf02OnDate = sDf02Imdg.getCollectionObjectsByPredicate(predicate); + log.debug("Found {} sdf09 records in range {} - {}", sdf09OnDate.size(), tradingDayStart, tradingDayEnd); + log.debug("Found {} sdf02 records in range {} - {}", sdf02OnDate.size(), tradingDayStart, tradingDayEnd); + + String message = sdf09OnDate.isEmpty() && sdf02OnDate.isEmpty() ? "Отсутствует пара ДФ-01/ДФ-57 и ДФ-08/ДФ-21" + : sdf09OnDate.isEmpty() ? "Отсутствует пара ДФ-08/ДФ-21" : sdf02OnDate.isEmpty() ? "Отсутствует пара ДФ-01/ДФ-57" : null; + if (message != null) { + final String destination = Consts.NOTIFICATION_NEW; + NotificationNewRequest request = new NotificationNewRequest(); + request.setObjectType("RGST"); + request.setPriority(Priority.HIGH.getKey()); + request.setComment(message); + log.debug("Send message to kafka \"{}\": {}", destination, LogFormatter.toStringWrapper(request)); + kafkaSender.sendRequestToQueue(destination, request); + } + } + +} diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java index 78172c922..aa5a7e74e 100644 --- a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java @@ -53,6 +53,7 @@ public enum Task implements IEnumKey { loadIssue_LOSC("LOSC"), // Загрузка инструментов dbfExport_OUTV("OUTV"), // для dbf-exporter сообщение на создание файла ДФ-54 sdf05WithCode9Final("FDFF"), //Формирование ДФ-05 с кодом 9 (финальный) + checkSdfPair("CHDF"),//Проверка наличия пары ДФ-01/ДФ-57 и ДФ-08/ДФ-21 ; private final String key; diff --git a/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/time/TimeUtil.java b/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/time/TimeUtil.java index d6efc9f5e..43624f183 100644 --- a/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/time/TimeUtil.java +++ b/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/time/TimeUtil.java @@ -3,6 +3,7 @@ package ru.spcex.platform.utils.time; import java.sql.Time; import java.time.*; import java.time.format.DateTimeFormatter; +import java.time.temporal.ChronoUnit; import java.util.Date; import java.util.Objects; @@ -79,4 +80,16 @@ public class TimeUtil { if (zoned == null) return null; return zoned.toLocalDateTime(); } + + /** + * Начало дня по текущей временной зоне. + * + * @param from + * @return + */ + public static Instant startOfDay(Instant from) { + ZonedDateTime zdtStart = from.atZone(ZoneId.systemDefault()); + zdtStart = zdtStart.truncatedTo(ChronoUnit.DAYS); + return zdtStart.toInstant(); + } }