---
AdmittedLiabilitiesRegisterService
CoveredLiabilitiesRegisterService
DepoBalanceRegisterService
ExcludeLiabilitiesRegisterService
refactored work with imdg and added some logs
This commit is contained in:
aalehin 2023-05-23 01:35:38 +03:00
parent 7918f78761
commit ee8ef5c24c
7 changed files with 125 additions and 45 deletions

View file

@ -14,13 +14,11 @@ import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.platform.enumeration.*;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import java.util.Collection;
import java.util.HashMap;
import java.util.List;
import java.util.Optional;
import java.util.*;
import static ru.spcex.platform.enumeration.Task.createRegistry_GRRT;
@ -54,17 +52,10 @@ public class AdmittedLiabilitiesRegisterService extends QueueConsumer implements
public void admittedLiabilitiesRegisterNew(BaseRequest<LauncherCommandRequest> userRequest) {
log.debug("LauncherCommandRequest received from {}", createRegistry_GRRT.topic()); //TODO FIX UP WHAT KIND OF TOPIC
List<Registry> registries = registryMap.getAllValues().stream().filter((x) -> {
return x.getRegistryDesignation().equalsIgnoreCase("O") &&
x.getRegistryCapacity().equalsIgnoreCase("P") &&
x.getRegistryCode().equalsIgnoreCase("T") &&
(x.getRegistryInstrumentType().equalsIgnoreCase("S")
|| x.getRegistryInstrumentType().equalsIgnoreCase("M"));
}
).toList();
String sqlConditionForRegistry = getSqlForRegistries();
List<Registry> registries = registryMap.getCollectionObjectsBySQL(sqlConditionForRegistry).stream().toList();
HashMap<Long, Long> sessionIdByCompanyId = new HashMap<>();
Collection<AdmittedLiabilitiesRegister> collection = admittedLiabilitiesRegisterMap.getAllValues();
Collection<AdmittedLiabilitiesRegister> collection = admittedLiabilitiesRegisterMap.getAllValues(); //there should be checked all values cause we can have a lot of registries
collection.forEach(x -> sessionIdByCompanyId.put(x.getCompanyId(), x.getSessionId()));
registries.forEach(x -> {
Long existingCompanySessionId = sessionIdByCompanyId.get(x.getCompanyId());
@ -77,8 +68,8 @@ public class AdmittedLiabilitiesRegisterService extends QueueConsumer implements
private void insertAdmittedLiabilitiesRegister(Registry registry) {
AdmittedLiabilitiesRegister admittedLiabilitiesRegister = new AdmittedLiabilitiesRegister();
Optional<Company> company = companyMap.getAllValues().stream().filter((x)-> x.getId() == 1L).findFirst();
Optional<Security> security = securityMap.getAllValues().stream().filter((x)-> x.getId().equals(registry.getSecurityId())).findFirst();
Optional<Company> company = companyMap.getCollectionObjectsByFieldValues(Map.of("id", "1")).stream().findFirst();
Optional<Security> security = securityMap.getCollectionObjectsByFieldValues(Map.of("id", registry.getSecurityId())).stream().findFirst();
String securitySymbol = security.map(Security::getSecuritySymbol).orElse(null);
String securityFullName = security.map(Security::getShortName).orElse(null);
String companyFullName = company.map(Company::getFullName).orElse(null);
@ -93,4 +84,16 @@ public class AdmittedLiabilitiesRegisterService extends QueueConsumer implements
admittedLiabilitiesRegister.setClearingDate(registry.getClearingDate());
admittedLiabilitiesRegisterMap.insert(admittedLiabilitiesRegister);
}
private String getSqlForRegistries(){
return String.format("registryDesignation = '%s' and " +
"registryInstrumentType in ('%s','%s') and " +
"registryCapacity = '%s' and " +
"registryCode = '%s'"
,
RegistryDesignation.O.getKey(),
RegistryInstrumentType.S.getKey(),
RegistryInstrumentType.M.getKey(),
RegistryCapacity.P.getKey(),
RegistryCode.T.getKey());
}
}

View file

@ -14,6 +14,7 @@ import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.platform.enumeration.*;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
@ -56,24 +57,15 @@ public class CoveredLiabilitiesRegisterService extends QueueConsumer implements
log.debug("LauncherCommandRequest received from {}", createRegistry_GRRT.topic()); //TODO FIX UP WHAT KIND OF TOPIC
HashMap<Long, Long> sessionIdByCompanyId = new HashMap<>();
List<Registry> registries = registryMap.getAllValues().stream().filter((x) -> {
return x.getRegistryDesignation().equalsIgnoreCase("O") &&
x.getRegistryCapacity().equalsIgnoreCase("P") &&
x.getRegistryCode().equalsIgnoreCase("T") &&
(x.getRegistryInstrumentType().equalsIgnoreCase("S")
|| x.getRegistryInstrumentType().equalsIgnoreCase("M")) &&
(x.getRegistryStatus().equalsIgnoreCase("OK") ||
x.getRegistryStatus().equalsIgnoreCase("PROC"));
}
String sqlConditionForRegistry = getSqlForRegistries();
).toList();
List<Registry> registries = registryMap.getCollectionObjectsBySQL(sqlConditionForRegistry).stream().toList();
Collection<CoveredLiabilitiesRegister> collection = coveredLiabilitiesRegisterMap.getAllValues();
collection.forEach(x -> sessionIdByCompanyId.put(x.getCompanyId(), x.getSessionId()));
registries.forEach(x -> {
Long existingCompanySessionId = sessionIdByCompanyId.get(x.getCompanyId());
if (existingCompanySessionId == null) {
insertCoveredLiabilitiesRegister(x);
}
});
log.debug("successfully processed");
@ -96,4 +88,19 @@ public class CoveredLiabilitiesRegisterService extends QueueConsumer implements
coveredLiabilitiesRegister.setClearingDate(registry.getClearingDate());
coveredLiabilitiesRegisterMap.insert(coveredLiabilitiesRegister);
}
private String getSqlForRegistries(){
return String.format("registryDesignation = '%s' and " +
"registryInstrumentType in ('%s','%s') and " +
"registryCapacity = '%s' and " +
"registryCode = '%s' and" +
"registryStatus in ('%s', '%s')"
,
RegistryDesignation.O.getKey(),
RegistryInstrumentType.S.getKey(),
RegistryInstrumentType.M.getKey(),
RegistryCapacity.P.getKey(),
RegistryCode.T.getKey(),
RegistryStatus.OK.getKey(),
RegistryStatus.PROC.getKey());
}
}

View file

@ -13,6 +13,7 @@ import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.platform.enumeration.*;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
@ -49,9 +50,11 @@ public class DepoBalanceRegisterService extends QueueConsumer implements Initial
public void depoBalanceRegisterNew(BaseRequest<LauncherCommandRequest> userRequest) {
log.debug("LauncherCommandRequest received from {}", createRegistry_GRRT.topic());
HashMap<Long, Long> sessionIdByCompanyId = new HashMap<>();
List<Registry> registries = registryMap.getAllValues().stream().filter((x) -> x.getAccountType().equalsIgnoreCase("DEPO")).toList();
String sqlForRegistries = getSqlForRegistries();
List<Registry> registries = registryMap.getCollectionObjectsBySQL(sqlForRegistries).stream().toList();
log.trace("Started searching depoBalanceRegister in register by companyId, sessionId...");
Collection<DepoBalanceRegister> collection = depoBalanceRegisterMap.getAllValues();
collection.forEach(x -> sessionIdByCompanyId.put(x.getCompanyId(), x.getSessionId()));
registries.forEach(x -> {
@ -60,10 +63,11 @@ public class DepoBalanceRegisterService extends QueueConsumer implements Initial
insertDepoBalanceRegister(x);
}
});
log.debug("successfully processed");
log.debug("Successfully processed");
}
private void insertDepoBalanceRegister(Registry registry) {
log.trace("Started generating DepoBalanceRegister entity...");
DepoBalanceRegister depoBalanceRegister = new DepoBalanceRegister();
depoBalanceRegister.setUpdated(Instant.now());
depoBalanceRegister.setCreated(Instant.now());
@ -73,5 +77,11 @@ public class DepoBalanceRegisterService extends QueueConsumer implements Initial
depoBalanceRegister.setQuantity(registry.getCloseBalance());
depoBalanceRegister.setSecuritySymbol(registry.getSecuritySymbol());
depoBalanceRegisterMap.insert(depoBalanceRegister);
log.debug("inserted successfully DepoBalanceRegister entity with id: {}", depoBalanceRegister.getId());
}
private String getSqlForRegistries(){
log.trace("Started generating sql predicate for registries...");
return String.format("accountType = '%s'",
AccountType.Depo.getKey());
}
}

View file

@ -2,22 +2,27 @@ package ru.spcex.clearing.registry.service;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.logging.log4j.util.StringMap;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Autowired;
import ru.clearing.classes.statics.data.company.CompanySymbols;
import ru.clearing.classes.statics.data.misc.Session;
import ru.clearing.classes.statics.data.register.ExcludeLiabilitiesRegister;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.clearing.classes.statics.data.security.MoneyMarketSecurity;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.platform.enumeration.*;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import java.util.Collection;
import java.util.HashMap;
import java.util.List;
import java.time.LocalDate;
import java.time.ZoneId;
import java.util.*;
import static ru.spcex.platform.enumeration.Task.createRegistry_GRRT;
@ -26,6 +31,9 @@ public class ExcludeLiabilitiesRegisterService extends QueueConsumer implements
private final Imdg<ExcludeLiabilitiesRegister> excludeLiabilitiesRegisterMap;
private final Imdg<Registry> registryMap;
private final Imdg<Session> sessionMap;
private final Imdg<ru.clearing.classes.statics.data.company.CompanySymbols> companySymbolMap;
private final Imdg<MoneyMarketSecurity> moneyMarketSecurityMap;
@Autowired
public ExcludeLiabilitiesRegisterService(Consumer<String, Object> kafkaQueue,
@ -34,6 +42,9 @@ public class ExcludeLiabilitiesRegisterService extends QueueConsumer implements
super(kafkaQueue, kafkaProducer);
this.excludeLiabilitiesRegisterMap = imdgProvider.getImdg(IMDGDistributedNames.Map_ExcludeLiabilitiesRegister, ExcludeLiabilitiesRegister.class);
this.registryMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
this.sessionMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class);
this.companySymbolMap = imdgProvider.getImdg(IMDGDistributedNames.Map_CompanySymbols, ru.clearing.classes.statics.data.company.CompanySymbols.class);
this.moneyMarketSecurityMap = imdgProvider.getImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, MoneyMarketSecurity.class);
}
@Override
@ -45,18 +56,11 @@ public class ExcludeLiabilitiesRegisterService extends QueueConsumer implements
}
public void excludeLiabilitiesRegisterNew(BaseRequest<LauncherCommandRequest> userRequest) {
public void excludeLiabilitiesRegisterNew(BaseRequest<LauncherCommandRequest> userRequest) {
log.debug("LauncherCommandRequest received from {}", createRegistry_GRRT.topic()); //TODO FIX UP WHAT KIND OF TOPIC
HashMap<Long, Long> sessionIdByCompanyId = new HashMap<>();
List<Registry> registries = registryMap.getAllValues().stream().filter((x) -> {
return x.getRegistryDesignation().equalsIgnoreCase("O") &&
x.getRegistryCapacity().equalsIgnoreCase("P") &&
x.getRegistryCode().equalsIgnoreCase("T") &&
(x.getRegistryInstrumentType().equalsIgnoreCase("S")
|| x.getRegistryInstrumentType().equalsIgnoreCase("M")) &&
x.getRegistryStatus().equalsIgnoreCase("FAIL");
}
).toList();
String sqlForRegistries = getSqlForRegistries();
List<Registry> registries = registryMap.getCollectionObjectsBySQL(sqlForRegistries).stream().toList();
Collection<ExcludeLiabilitiesRegister> collection = excludeLiabilitiesRegisterMap.getAllValues();
collection.forEach(x -> sessionIdByCompanyId.put(x.getCompanyId(), x.getSessionId()));
registries.forEach(x -> {
@ -70,7 +74,40 @@ public class ExcludeLiabilitiesRegisterService extends QueueConsumer implements
private void insertExcludeLiabilitiesRegister(Registry registry) {
ExcludeLiabilitiesRegister excludeLiabilitiesRegister = new ExcludeLiabilitiesRegister();
//TODO MAPPING
//hz communicating
Session sessionBySessionId = sessionMap.getSingleObjectByFieldValues(Map.of("clearingDate", registry.getSessionId()));
LocalDate validToDate = LocalDate.ofInstant(sessionBySessionId.getUpdated(), ZoneId.systemDefault());
Map<String, ? extends Comparable<?>> innPredicates = Map.of("companySymbol", "INN","companyId", registry.getCompanyId());
ru.clearing.classes.statics.data.company.CompanySymbols companySymbol = companySymbolMap.getSingleObjectByFieldValues(innPredicates);
MoneyMarketSecurity moneyMarketSecurity = moneyMarketSecurityMap.getSingleObjectByFieldValues(Map.of("securityId", registry.getSecurityId()));
excludeLiabilitiesRegister.setSessionId(registry.getSessionId());
excludeLiabilitiesRegister.setValidFromDate(sessionBySessionId.getClearingDate());
excludeLiabilitiesRegister.setValidToDate(validToDate); //TODO CLARIFY IS IN THIS FIELD SHOULD LOCAL DATE OR NO?
excludeLiabilitiesRegister.setCompanyId(registry.getCompanyId());
excludeLiabilitiesRegister.setCompanyFullName(""); //TODO CLARIFY HOW IT SHOULD BE FILLED
excludeLiabilitiesRegister.setInn(companySymbol.getCompanySymbolValue());
excludeLiabilitiesRegister.setRegistryCode(registry.getRegistryCode());
excludeLiabilitiesRegister.setAccount(registry.getAccount());
excludeLiabilitiesRegister.setRegistryStatus(registry.getRegistryStatus());
excludeLiabilitiesRegister.setCurrency(moneyMarketSecurity.getNominalCurrency()); // TODO IS moneyMarketSecurity.currency NOMINAL CURRENCY
excludeLiabilitiesRegister.setSumLiabilities(registry.getBalance());
excludeLiabilitiesRegister.setSettlementDate(registry.getSettlementDate());
excludeLiabilitiesRegisterMap.insert(excludeLiabilitiesRegister);
}
private String getSqlForRegistries() {
return String.format("registryDesignation = '%s' and " +
"registryInstrumentType in ('%s','%s') and " +
"registryCapacity = '%s' and " +
"registryCode = '%s' and " +
"registryStatus = '%s'"
,
RegistryDesignation.O.getKey(),
RegistryInstrumentType.S.getKey(),
RegistryInstrumentType.M.getKey(),
RegistryCapacity.P.getKey(),
RegistryCode.T.getKey(),
RegistryStatus.FAIL.getKey());
}
}

View file

@ -9,6 +9,7 @@ public enum RegistryCapacity implements IEnumKey {
Z("Z"),
C("C"),
E("E"),
P("P")
;
private final String key;

View file

@ -0,0 +1,19 @@
package ru.spcex.platform.enumeration;
import ru.spcex.platform.utils.enumeration.IEnumKey;
public enum RegistryCode implements IEnumKey {
T("T");
//TODO где взять дополнительные значения для registry code?
private final String key;
RegistryCode(String key) {
this.key = key;
}
@Override
public String getKey() {
return key;
}
}

View file

@ -3,7 +3,10 @@ package ru.spcex.platform.enumeration;
import ru.spcex.platform.utils.enumeration.IEnumKey;
public enum RegistryStatus implements IEnumKey {
OK("OK"), PROC("PROC"), NACK("NACK"),
OK("OK"),
PROC("PROC"),
NACK("NACK"),
FAIL("FAIL")
;
private final String key;