This commit is contained in:
parent
5fd85852bb
commit
723f177c25
3 changed files with 102 additions and 0 deletions
|
|
@ -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<SDf02> sDf02Imdg;
|
||||
private final Imdg<SDf09> sDf09Imdg;
|
||||
private final ImdgProvider imdgProvider;
|
||||
private final KafkaSender kafkaSender;
|
||||
|
||||
@Autowired
|
||||
public SdfPairChecker(ImdgProvider imdgProvider,
|
||||
KafkaSender kafkaSender,
|
||||
Producer<String, Object> kafkaProducer,
|
||||
Consumer<String, Object> 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<LauncherCommandRequest> 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<SDf09> sdf09OnDate = sDf09Imdg.getCollectionObjectsByPredicate(predicate);
|
||||
Collection<SDf02> 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);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue