http://jira.mfd.msk:8088/browse/CLS-781 проверка 08 21 для ЦК
This commit is contained in:
parent
e007a42226
commit
85d9fc6154
3 changed files with 29 additions and 8 deletions
|
|
@ -14,6 +14,15 @@ public class UtilityServiceSettings {
|
||||||
private HazelcastClientParams hazelcast;
|
private HazelcastClientParams hazelcast;
|
||||||
private KafkaConsumerSettings kafkaConsumer;
|
private KafkaConsumerSettings kafkaConsumer;
|
||||||
private KafkaProducerSettings kafkaProducer;
|
private KafkaProducerSettings kafkaProducer;
|
||||||
|
private Boolean includeCk = false;
|
||||||
|
|
||||||
|
public Boolean getIncludeCk() {
|
||||||
|
return includeCk;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setIncludeCk(Boolean includeCk) {
|
||||||
|
this.includeCk = includeCk;
|
||||||
|
}
|
||||||
|
|
||||||
public HazelcastClientParams getHazelcast() {
|
public HazelcastClientParams getHazelcast() {
|
||||||
return hazelcast;
|
return hazelcast;
|
||||||
|
|
|
||||||
|
|
@ -1,5 +1,8 @@
|
||||||
package ru.spcex.clearing.utility.service;
|
package ru.spcex.clearing.utility.service;
|
||||||
|
|
||||||
|
import java.time.Instant;
|
||||||
|
import java.time.temporal.ChronoUnit;
|
||||||
|
import java.util.Collection;
|
||||||
import org.apache.kafka.clients.consumer.Consumer;
|
import org.apache.kafka.clients.consumer.Consumer;
|
||||||
import org.apache.kafka.clients.producer.Producer;
|
import org.apache.kafka.clients.producer.Producer;
|
||||||
import org.slf4j.Logger;
|
import org.slf4j.Logger;
|
||||||
|
|
@ -17,6 +20,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationNew
|
||||||
import ru.spcex.clearing.platform.messaging.serialization.LogFormatter;
|
import ru.spcex.clearing.platform.messaging.serialization.LogFormatter;
|
||||||
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.clearing.utility.config.settings.UtilityServiceSettings;
|
||||||
import ru.spcex.platform.enumeration.Priority;
|
import ru.spcex.platform.enumeration.Priority;
|
||||||
import ru.spcex.platform.enumeration.Task;
|
import ru.spcex.platform.enumeration.Task;
|
||||||
import ru.spcex.platform.imdg.api.Imdg;
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
|
|
@ -26,10 +30,6 @@ import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
|
||||||
import ru.spcex.platform.imdg.iml.hazelcast.adapter.predicate.ImdgPredicateBuilderHazelcast;
|
import ru.spcex.platform.imdg.iml.hazelcast.adapter.predicate.ImdgPredicateBuilderHazelcast;
|
||||||
import ru.spcex.platform.utils.time.TimeUtil;
|
import ru.spcex.platform.utils.time.TimeUtil;
|
||||||
|
|
||||||
import java.time.Instant;
|
|
||||||
import java.time.temporal.ChronoUnit;
|
|
||||||
import java.util.Collection;
|
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
public class SdfPairChecker extends QueueConsumer implements InitializingBean {
|
public class SdfPairChecker extends QueueConsumer implements InitializingBean {
|
||||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
|
@ -37,17 +37,20 @@ public class SdfPairChecker extends QueueConsumer implements InitializingBean {
|
||||||
private final Imdg<SDf09> sDf09Imdg;
|
private final Imdg<SDf09> sDf09Imdg;
|
||||||
private final ImdgProvider imdgProvider;
|
private final ImdgProvider imdgProvider;
|
||||||
private final KafkaSender kafkaSender;
|
private final KafkaSender kafkaSender;
|
||||||
|
private final boolean includeCk;
|
||||||
|
|
||||||
@Autowired
|
@Autowired
|
||||||
public SdfPairChecker(ImdgProvider imdgProvider,
|
public SdfPairChecker(ImdgProvider imdgProvider,
|
||||||
KafkaSender kafkaSender,
|
KafkaSender kafkaSender,
|
||||||
Producer<String, Object> kafkaProducer,
|
Producer<String, Object> kafkaProducer,
|
||||||
Consumer<String, Object> kafkaQueue) {
|
Consumer<String, Object> kafkaQueue,
|
||||||
|
UtilityServiceSettings settings) {
|
||||||
super(kafkaQueue, kafkaProducer);
|
super(kafkaQueue, kafkaProducer);
|
||||||
this.imdgProvider = imdgProvider;
|
this.imdgProvider = imdgProvider;
|
||||||
this.kafkaSender = kafkaSender;
|
this.kafkaSender = kafkaSender;
|
||||||
this.sDf02Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf02, SDf02.class);
|
this.sDf02Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf02, SDf02.class);
|
||||||
this.sDf09Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf09, SDf09.class);
|
this.sDf09Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf09, SDf09.class);
|
||||||
|
this.includeCk = settings.getIncludeCk();
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|
@ -72,8 +75,15 @@ public class SdfPairChecker extends QueueConsumer implements InitializingBean {
|
||||||
log.debug("Found {} sdf09 records in range {} - {}", sdf09OnDate.size(), tradingDayStart, tradingDayEnd);
|
log.debug("Found {} sdf09 records in range {} - {}", sdf09OnDate.size(), tradingDayStart, tradingDayEnd);
|
||||||
log.debug("Found {} sdf02 records in range {} - {}", sdf02OnDate.size(), tradingDayStart, tradingDayEnd);
|
log.debug("Found {} sdf02 records in range {} - {}", sdf02OnDate.size(), tradingDayStart, tradingDayEnd);
|
||||||
|
|
||||||
String message = sdf09OnDate.isEmpty() && sdf02OnDate.isEmpty() ? "Отсутствуют пары ДФ-01/ДФ-57 и ДФ-08/ДФ-21"
|
String message;
|
||||||
: sdf09OnDate.isEmpty() ? "Отсутствует пара ДФ-08/ДФ-21" : sdf02OnDate.isEmpty() ? "Отсутствует пара ДФ-01/ДФ-57" : null;
|
{
|
||||||
|
boolean sdf09Err = !includeCk && sdf09OnDate.isEmpty();
|
||||||
|
boolean sdf02Err = sdf02OnDate.isEmpty();
|
||||||
|
message = sdf09Err && sdf02Err ? "Отсутствуют пары ДФ-01/ДФ-57 и ДФ-08/ДФ-21"
|
||||||
|
: sdf09Err ? "Отсутствует пара ДФ-08/ДФ-21"
|
||||||
|
: sdf02Err ? "Отсутствует пара ДФ-01/ДФ-57"
|
||||||
|
: null;
|
||||||
|
}
|
||||||
if (message != null) {
|
if (message != null) {
|
||||||
final String destination = Consts.NOTIFICATION_NEW;
|
final String destination = Consts.NOTIFICATION_NEW;
|
||||||
NotificationNewRequest request = new NotificationNewRequest();
|
NotificationNewRequest request = new NotificationNewRequest();
|
||||||
|
|
|
||||||
|
|
@ -17,4 +17,6 @@ utility-service.kafka-producer.acks=all
|
||||||
utility-service.kafka-producer.retries=0
|
utility-service.kafka-producer.retries=0
|
||||||
utility-service.kafka-producer.batch-size=16384
|
utility-service.kafka-producer.batch-size=16384
|
||||||
utility-service.kafka-producer.linger-ms=1
|
utility-service.kafka-producer.linger-ms=1
|
||||||
utility-service.kafka-producer.buffer-memory=33554432
|
utility-service.kafka-producer.buffer-memory=33554432
|
||||||
|
|
||||||
|
utility-service.include-ck=false
|
||||||
Loading…
Add table
Reference in a new issue