utility-service http://jira.mfd.msk:8088/browse/CLS-839 сверка AM_T,AM_F по GVTF

This commit is contained in:
AKurakin 2025-05-05 18:29:51 +03:00
parent b6d6db0499
commit 40f248ad7c
4 changed files with 401 additions and 1 deletions

View file

@ -0,0 +1,223 @@
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.registry.Registry;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
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.*;
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.api.predicate.specific.RegistryCodeSqlBuilder;
import ru.spcex.platform.imdg.iml.hazelcast.adapter.predicate.ImdgPredicateBuilderHazelcast;
import ru.spcex.platform.utils.collection.Pair;
import ru.spcex.platform.utils.number.BigDecimalUtil;
import java.math.BigDecimal;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.*;
import java.util.stream.Collectors;
@Service
public class AmatAmafChecker extends QueueConsumer implements InitializingBean {
private final Logger log = LoggerFactory.getLogger(getClass());
static final DateTimeFormatter MESSAGE_TIME_FORMATTER = DateTimeFormatter.ofPattern("yyyy.MM.dd HH:mm:ss");
private final Imdg<Registry> registryImdg;
private final Imdg<TradingClearingRegistry> tradingClearingRegistryImdg;
private final ImdgProvider imdgProvider;
private final KafkaSender kafkaSender;
@Autowired
public AmatAmafChecker(ImdgProvider imdgProvider,
KafkaSender kafkaSender,
Producer<String, Object> kafkaProducer,
Consumer<String, Object> kafkaQueue) {
super(kafkaQueue, kafkaProducer);
this.imdgProvider = imdgProvider;
this.kafkaSender = kafkaSender;
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
this.tradingClearingRegistryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
}
@Override
public void afterPropertiesSet() throws Exception {
callback(LauncherCommandRequest.class)
.setConsumer(this::checkOnDate)
.forDestination(Task.GVTF.topic(), callbacks::put);
init();
}
public void checkOnDate(BaseRequest<LauncherCommandRequest> gvtfEvent) {
List<Long> specialId = null;
if (gvtfEvent.getRequestPayload() != null && gvtfEvent.getRequestPayload().getIds() != null) {
specialId = gvtfEvent.getRequestPayload().getIds();
}
log.debug("Check AM*T and AM*F (parameter id's = {})", specialId);
ImdgPredicateBuilder builder = ImdgPredicateBuilderHazelcast.instance();
Collection<Registry> checkRegistries;
final LocalDateTime checkTime = LocalDateTime.now();
ImdgPredicate AM_T_AM_F_predicate = RegistryCodeSqlBuilder.getInstance(
RegistryTradingParams.AM_T,
RegistryTradingParams.AM_F)
.buildPredicate(builder);
if (specialId == null) {
long timer = System.currentTimeMillis();
Collection<TradingClearingRegistry> activeTcr = tradingClearingRegistryImdg.getCollectionObjectsByPredicate(
builder.equals("status", ServiceStatus.Active.getKey())
);
List<String> tcrActiveCode = activeTcr.stream()
.map(tcr -> tcr.getCode())
.filter(Objects::nonNull)
.collect(Collectors.toList());
log.debug("Found {} code's from active TradingClearingRegistry", tcrActiveCode.size());
ImdgPredicate predicate = builder.and(
AM_T_AM_F_predicate,
builder.equals("accountType", AccountType.Info.getKey()),
builder.in("tradingClearingRegistry", tcrActiveCode.toArray(new String[0]))
);
checkRegistries = registryImdg.getCollectionObjectsByPredicate(predicate);
timer = System.currentTimeMillis() - timer;
log.debug("Select {} Registry by query: {}. Time {} ms.", checkRegistries.size(), predicate, timer);
} else {
log.debug("Selecting {} AM_T and by tradingClearingRegistry and securitySymbol AM_F for his.", specialId);
long timer = System.currentTimeMillis();
checkRegistries = new ArrayList<>(specialId.size() * 2);
for (Long am_tId : specialId) {
if (am_tId == null) { // never
log.warn("Received argument ID was null");
continue;
}
Registry am_tRegistry = registryImdg.getSingleObjectByID(am_tId);
if (am_tRegistry == null) {
log.warn("Registry not found by ID={}", am_tId);
} else {
checkRegistries.add(am_tRegistry);
ImdgPredicate predicate = builder.and(
AM_T_AM_F_predicate,
builder.equals("accountType", AccountType.Info.getKey()),
builder.equals("tradingClearingRegistry", am_tRegistry.getTradingClearingRegistry()),
builder.equals("securitySymbol", am_tRegistry.getSecuritySymbol()),
builder.not(builder.equals("id", am_tRegistry.getId())) // counter-registry
);
Collection<Registry> otherRegistry = registryImdg.getCollectionObjectsByPredicate(predicate);
log.trace("For Registry.id={} found {} other Registry's by query {}", am_tId, otherRegistry.size(), predicate);
for (Registry otherReg : otherRegistry) {
if (!am_tId.equals(otherReg.getId()))
checkRegistries.add(am_tRegistry);
}
}
}
timer = System.currentTimeMillis() - timer;
log.debug("Select {} Registry and hi's pairs by ID collection. Time {} ms.", checkRegistries.size(), timer);
}
// Сверка Registry
List<String> mismatch = checkRegisters(checkRegistries);
if (mismatch.isEmpty()) {
log.debug("Check registry result: OK");
sendMessage(makeSuccessMessage(checkTime), Priority.LOW);
} else {
log.debug("Check registry result: not match");
sendMessage(makeErrorMessage(checkTime, mismatch), Priority.HIGH);
}
}
String makeSuccessMessage(LocalDateTime checkTime) {
return String.format("""
Проведена сверка по регистрам AMAT/AMAF.
%s.
Результат сверки: расхождения отсутствуют.""",
MESSAGE_TIME_FORMATTER.format(checkTime));
}
String makeErrorMessage(LocalDateTime checkTime, List<String> mismatch) {
assert !mismatch.isEmpty();
String msgItems = mismatch.stream().collect(Collectors.joining("; "));
return String.format("""
Проведена сверка по регистрам AMAT/AMAF.
%s.
Результат сверки: %s.""",
MESSAGE_TIME_FORMATTER.format(checkTime), msgItems);
}
protected List<String> checkRegisters(Collection<Registry> checkRegistries) {
List<String> resultMismatch = new ArrayList<>();
Map<Pair<String, String>, List<Registry>> regGroup = checkRegistries.stream().collect(Collectors.groupingBy(
r -> new Pair<>(r.getTradingClearingRegistry(), r.getSecuritySymbol())
));
log.trace("Make {} group of TradingClearingRegistry+SecuritySymbol from {} Registry", regGroup.size(), checkRegistries.size());
int equalsCount = 0;
int mismatchCount = 0;
for (Map.Entry<Pair<String, String>, List<Registry>> group : regGroup.entrySet()) {
BigDecimal am_tSumm = BigDecimal.ZERO;
BigDecimal am_fSumm = BigDecimal.ZERO;
for (Registry r : group.getValue()) {
if (match(RegistryTradingParams.AM_T, r)) {
am_tSumm = am_tSumm.add(BigDecimalUtil.safeBD(r.getBalance()));
}
if (match(RegistryTradingParams.AM_F, r)) {
am_fSumm = am_fSumm.add(BigDecimalUtil.safeBD(r.getBalance()));
}
}
log.trace("Group {} size of {} item: AM_T summ={}, AM_F summ={}",
group.getKey(), group.getValue().size(), am_tSumm, am_fSumm);
boolean balanceOk = am_tSumm.compareTo(am_fSumm) == 0;
if (balanceOk) {
equalsCount++;
} else {
mismatchCount++;
String describe = String.format("по ТКР %s выявлено расхождение на регистрах AMAT/AMAF по %s на сумму %s",
group.getKey().getFirst(), // TradingClearingRegistry
group.getKey().getSecond(), // SecuritySymbol
am_tSumm.subtract(am_fSumm)
);
resultMismatch.add(describe);
}
}
log.debug("Check result: equalsCount = {}, mismatchCount = {}", equalsCount, mismatchCount);
return resultMismatch;
}
protected void sendMessage(String message, Priority priority) {
final String destination = Consts.NOTIFICATION_NEW;
NotificationNewRequest request = new NotificationNewRequest();
request.setObjectType(ObjectType.rgst.getKey());
request.setPriority(priority.getKey());
request.setComment(message);
log.debug("Send message to kafka \"{}\": {}", destination, LogFormatter.toStringWrapper(request));
kafkaSender.sendRequestToQueue(destination, request);
}
protected boolean match(RegistryTradingParams params, Registry reg) {
if (params.registryDesignation() != null && !params.registryDesignation().equalsByKey(reg.getRegistryDesignation()))
return false;
if (params.registryCapacity() != null && !params.registryCapacity().equalsByKey(reg.getRegistryCapacity()))
return false;
if (params.registryInstrumentType() != null && !params.registryInstrumentType().equalsByKey(reg.getRegistryInstrumentType()))
return false;
if (params.registryUnit() != null && !params.registryUnit().equalsByKey(reg.getRegistryUnit()))
return false;
return true;
}
}

View file

@ -0,0 +1,147 @@
package ru.spcex.clearing.utility.service;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.producer.Producer;
import org.junit.jupiter.api.Test;
import org.mockito.Mockito;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.platform.enumeration.RegistryDesignation;
import ru.spcex.platform.enumeration.RegistryInstrumentType;
import ru.spcex.platform.enumeration.RegistryTradingParams;
import ru.spcex.platform.enumeration.RegistryUnit;
import ru.spcex.platform.imdg.api.ImdgProvider;
import java.math.BigDecimal;
import java.time.LocalDateTime;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
import static org.junit.jupiter.api.Assertions.*;
class AmatAmafCheckerTest {
AmatAmafChecker instance() {
ImdgProvider imdgProvider = Mockito.mock(ImdgProvider.class);
KafkaSender kafkaSender = Mockito.mock(KafkaSender.class);
Producer<String, Object> kafkaProducer = Mockito.mock(Producer.class);
Consumer<String, Object> kafkaQueue = Mockito.mock(Consumer.class);
return new AmatAmafChecker(imdgProvider, kafkaSender, kafkaProducer, kafkaQueue);
}
@Test
void checkOnDate() {
AmatAmafChecker checker = instance();
List<Registry> regSrc = new ArrayList<>();
List<String> tkr = Arrays.asList("0004CAV00003", "0482CAT00001", "0095CAV00002");
List<String> securitySymbol = Arrays.asList("RUB", "CNY");
int i = 1;
for (String tcrCode : tkr)
for (String secSymbol : securitySymbol) {
Registry reg1 = new Registry();
reg1.setId(1L);
reg1.setRegistryDesignation("A");
reg1.setRegistryInstrumentType("M");
reg1.setRegistryCapacity("A");
reg1.setRegistryUnit("T");
reg1.setRegistryCode("AMAT"); // RegistryUtil.clearingCode(reg1)
reg1.setTradingClearingRegistry(tcrCode);
reg1.setSecuritySymbol(secSymbol);
reg1.setBalance(BigDecimal.valueOf(100 * i));
if ("RUB".equals(secSymbol))
reg1.setBalanceRub(reg1.getBalance());
else
reg1.setBalanceRub(reg1.getBalance().multiply(BigDecimal.valueOf(0.1)));
regSrc.add(reg1);
Registry reg2 = new Registry();
reg2.setId(1L);
reg2.setRegistryDesignation("A");
reg2.setRegistryInstrumentType("M");
reg2.setRegistryCapacity("A");
reg2.setRegistryUnit("F");
reg2.setRegistryCode("AMAF");
reg2.setTradingClearingRegistry(tcrCode);
reg2.setSecuritySymbol(secSymbol);
reg2.setBalance(reg1.getBalance());
reg2.setBalanceRub(reg1.getBalanceRub());
regSrc.add(reg2);
}
// All OK
List<String> errors = checker.checkRegisters(regSrc);
assertEquals(0, errors.size());
// Error in data
regSrc.get(0).setBalance(regSrc.get(0).getBalance().add(BigDecimal.valueOf(100)));
regSrc.get(4).setBalance(BigDecimal.valueOf(10.25));
errors = checker.checkRegisters(regSrc);
assertEquals(2, errors.size());
assertEquals("[по ТКР 0004CAV00003 выявлено расхождение на регистрах AMAT/AMAF по RUB на сумму 100, по ТКР 0482CAT00001 выявлено расхождение на регистрах AMAT/AMAF по RUB на сумму -89.75]",
errors.toString());
// next test makeErrorMessage:
assertEquals("Проведена сверка по регистрам AMAT/AMAF.\n" +
"2025.05.05 12:37:56.\n" +
"Результат сверки: по ТКР 0004CAV00003 выявлено расхождение на регистрах AMAT/AMAF по RUB на сумму 100; по ТКР 0482CAT00001 выявлено расхождение на регистрах AMAT/AMAF по RUB на сумму -89.75.",
checker.makeErrorMessage(
LocalDateTime.of(2025, 5, 5, 12, 37, 56), errors));
}
@Test
void makeSuccessMessage() {
AmatAmafChecker checker = instance();
assertEquals("Проведена сверка по регистрам AMAT/AMAF.\n" +
"2025.05.05 12:37:56.\n" +
"Результат сверки: расхождения отсутствуют.",
checker.makeSuccessMessage(LocalDateTime.of(2025, 5, 5, 12, 37, 56)));
}
@Test
void makeErrorMessage() {
AmatAmafChecker checker = instance();
assertEquals("Проведена сверка по регистрам AMAT/AMAF.\n" +
"2025.05.05 12:37:56.\n" +
"Результат сверки: по ТКР 0004CAV00003 выявлено расхождение на регистрах AMAT/AMAF по RUB на сумму 100; по ТКР 0482CAT00001 выявлено расхождение на регистрах AMAT/AMAF по RUB на сумму -89.75.",
checker.makeErrorMessage(
LocalDateTime.of(2025, 5, 5, 12, 37, 56),
Arrays.asList(
"по ТКР 0004CAV00003 выявлено расхождение на регистрах AMAT/AMAF по RUB на сумму 100",
"по ТКР 0482CAT00001 выявлено расхождение на регистрах AMAT/AMAF по RUB на сумму -89.75"
)
));
}
@Test
void match() {
AmatAmafChecker checker = instance();
Registry reg1 = new Registry();
reg1.setId(1L);
reg1.setRegistryDesignation("A");
reg1.setRegistryInstrumentType("M");
reg1.setRegistryCapacity("A");
reg1.setRegistryUnit("T");
reg1.setRegistryCode("AMAT"); // RegistryUtil.clearingCode(reg1)
Registry reg2 = new Registry();
reg2.setId(1L);
reg2.setRegistryDesignation("A");
reg2.setRegistryInstrumentType("M");
reg2.setRegistryCapacity("A");
reg2.setRegistryUnit("F");
reg2.setRegistryCode("AMAF");
assertTrue(checker.match(RegistryTradingParams.AM_T, reg1));
assertTrue(checker.match(RegistryTradingParams.AM_F, reg2));
assertTrue(checker.match(RegistryTradingParams.AM__, reg1));
assertTrue(checker.match(RegistryTradingParams.AM__, reg2));
assertFalse(checker.match(RegistryTradingParams.AS__, reg1));
assertFalse(checker.match(RegistryTradingParams.AM_F, reg1));
assertFalse(checker.match(RegistryTradingParams.AM_T, reg2));
}
}

View file

@ -1,6 +1,5 @@
package ru.spcex.clearing.platform.messaging.domain.cud.schedule;
import com.fasterxml.jackson.annotation.JsonFormat;
import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.databind.annotation.JsonDeserialize;
import com.fasterxml.jackson.databind.annotation.JsonSerialize;
@ -12,6 +11,7 @@ import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalTimeSeria
import java.math.BigDecimal;
import java.time.LocalDate;
import java.time.LocalTime;
import java.util.List;
public class LauncherCommandRequest {
@ -66,6 +66,9 @@ public class LauncherCommandRequest {
@JsonProperty
public LocalDate toDate;
@JsonProperty
private List<Long> ids;
public Long getUserId() {
return userId;
}
@ -233,4 +236,12 @@ public class LauncherCommandRequest {
public void setToDate(LocalDate toDate) {
this.toDate = toDate;
}
public List<Long> getIds() {
return ids;
}
public void setIds(List<Long> ids) {
this.ids = ids;
}
}

View file

@ -1,5 +1,6 @@
package ru.spcex.platform.utils.collection;
import java.util.Objects;
import java.util.function.Consumer;
public class Pair<T1, T2> {
@ -40,4 +41,22 @@ public class Pair<T1, T2> {
secondConsumer.accept(second);
}
}
@Override
public boolean equals(Object o) {
if (this == o) return true;
if (o == null || getClass() != o.getClass()) return false;
Pair<?, ?> pair = (Pair<?, ?>) o;
return Objects.equals(first, pair.first) && Objects.equals(second, pair.second);
}
@Override
public int hashCode() {
return (first == null ? 0 : first.hashCode()) * 31 ^ (second == null ? 0 : second.hashCode());
}
@Override
public String toString() {
return "Pair{" + first + "; " + second + "}";
}
}