registry-service http://jira.mfd.msk:8088/browse/CLS-280 предварительная очистка мап перед вставкой; улучшил логи

This commit is contained in:
AKurakin 2023-09-29 12:25:27 +03:00
parent 36b0381294
commit e278bc1bb8
11 changed files with 181 additions and 37 deletions

View file

@ -17,6 +17,7 @@ import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.registry.AdmittedLiabilitiesRegisterNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.clearing.registry.util.PreClearMap;
import ru.spcex.platform.enumeration.RegistryDesignation;
import ru.spcex.platform.enumeration.RegistryInstrumentType;
import ru.spcex.platform.enumeration.RegistryStatus;
@ -24,10 +25,8 @@ import ru.spcex.platform.enumeration.RegistryUnit;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.specific.SecuritySelector;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.validation.IValidator;
import java.text.MessageFormat;
import java.time.Instant;
import java.util.Collection;
import java.util.Map;
@ -45,6 +44,7 @@ public class AdmittedLiabilitiesRegisterService extends QueueConsumer implements
private final Imdg<Company> companyMap;
private final SecuritySelector<Security> securitySelector;
private final Function<Map<String, ?>, IValidator> fieldValuesValidator;
private final PreClearMap<AdmittedLiabilitiesRegister> preClearMap;
@Autowired
public AdmittedLiabilitiesRegisterService(Consumer<String, Object> kafkaQueue,
@ -57,6 +57,7 @@ public class AdmittedLiabilitiesRegisterService extends QueueConsumer implements
this.companyMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
this.securitySelector = new SecuritySelector<>(imdgProvider, Security.class);
this.fieldValuesValidator = fieldValuesValidator;
preClearMap = PreClearMap.instanceForLocalDateField(admittedLiabilitiesRegisterMap, "clearingDate");
}
@Override
@ -78,6 +79,8 @@ public class AdmittedLiabilitiesRegisterService extends QueueConsumer implements
log.debug("No registry with such conditions {}", sqlConditionForRegistry);
return;
}
log.debug("Select {} Registry by query {}", registries.size(), registries);
preClearMap.preClearMap();
registries.forEach((registry -> {
if (isDuplicateInMap(registry).isEmpty()) {
insertAdmittedLiabilitiesRegister(registry);
@ -131,20 +134,7 @@ public class AdmittedLiabilitiesRegisterService extends QueueConsumer implements
@Override
public Optional<AdmittedLiabilitiesRegister> isDuplicateInMap(Registry entity) {
return Optional.empty();
//todo для всех isDuplicateInMap: сказали надо по критерию чистить мапу перед вставкой записей, критерий уточнят.
/*
todo вечер 28.09.2023
Вот по какому полю ориентируемся (название реестра - наименование поля)
depoBalanceRegister - createdAt
moneyBalanceRegister - createdAt
admittedLiabilitiesRegister - clearingDate
coveredLiabilitiesRegister - clearingDate
moneyPaymentInstructionRegister - clearingDate
depoPaymentInstructionRegister - createdAt
excludeLiabilitiesRegister - createdAt
liabilitiesRegister - createdAt
executionRegister - clearingDate
*/
//для всех isDuplicateInMap: пока не используется, надо по критерию чистить мапу перед вставкой записей, preClearMap.
// Map<String, ? extends Comparable<?>> fieldValues = Map.of(
// "companyId", entity.getRegistryCode(),
// "sessionId", entity.getSessionId(),

View file

@ -17,6 +17,7 @@ import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.registry.CoveredLiabilitiesRegisterNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.clearing.registry.util.PreClearMap;
import ru.spcex.platform.enumeration.RegistryDesignation;
import ru.spcex.platform.enumeration.RegistryInstrumentType;
import ru.spcex.platform.enumeration.RegistryStatus;
@ -24,10 +25,8 @@ import ru.spcex.platform.enumeration.RegistryUnit;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.specific.SecuritySelector;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.validation.IValidator;
import java.text.MessageFormat;
import java.time.Instant;
import java.util.Collection;
import java.util.Map;
@ -45,6 +44,7 @@ public class CoveredLiabilitiesRegisterService extends QueueConsumer implements
private final Imdg<Company> companyMap;
private final SecuritySelector<Security> securitySelector;
private final Function<Map<String, ?>, IValidator> fieldValuesValidator;
private final PreClearMap<CoveredLiabilitiesRegister> preClearMap;
@Autowired
public CoveredLiabilitiesRegisterService(Consumer<String, Object> kafkaQueue,
@ -57,6 +57,7 @@ public class CoveredLiabilitiesRegisterService extends QueueConsumer implements
this.companyMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
this.securitySelector = new SecuritySelector<>(imdgProvider, Security.class);
this.fieldValuesValidator = fieldValuesValidator;
preClearMap = PreClearMap.instanceForLocalDateField(coveredLiabilitiesRegisterMap, "clearingDate");
}
@Override
@ -79,6 +80,8 @@ public class CoveredLiabilitiesRegisterService extends QueueConsumer implements
log.debug("No registry with such conditions {}", sqlConditionForRegistry);
return;
}
log.debug("Select {} Registry by query {}", registries.size(), sqlConditionForRegistry);
preClearMap.preClearMap();
registries.forEach((registry -> {
if (isDuplicateInMap(registry).isEmpty()) {
insertCoveredLiabilitiesRegister(registry);

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.clearing.registry.util.PreClearMap;
import ru.spcex.platform.enumeration.AccountType;
import ru.spcex.platform.enumeration.RegistryTradingParams;
import ru.spcex.platform.imdg.api.Imdg;
@ -20,10 +21,8 @@ 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.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.validation.IValidator;
import java.text.MessageFormat;
import java.time.Instant;
import java.util.Collection;
import java.util.Map;
@ -39,6 +38,7 @@ public class DepoBalanceRegisterService extends QueueConsumer implements Initial
private final Imdg<DepoBalanceRegister> depoBalanceRegisterMap;
private final Imdg<Registry> registryMap;
private final Function<Map<String, ?>, IValidator> fieldValuesValidator;
private final PreClearMap<DepoBalanceRegister> preClearMap;
@Autowired
public DepoBalanceRegisterService(Consumer<String, Object> kafkaQueue,
@ -49,6 +49,7 @@ public class DepoBalanceRegisterService extends QueueConsumer implements Initial
this.depoBalanceRegisterMap = imdgProvider.getImdg(IMDGDistributedNames.Map_DepoBalanceRegister, DepoBalanceRegister.class);
this.registryMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
this.fieldValuesValidator = fieldValuesValidator;
preClearMap = PreClearMap.instanceForInstantField(depoBalanceRegisterMap, "created");
}
@Override
@ -68,6 +69,8 @@ public class DepoBalanceRegisterService extends QueueConsumer implements Initial
log.debug("No registry with such conditions '{}'", prdctForRegistries);
return;
}
log.debug("Select {} Registry by query {}", registries.size(), prdctForRegistries);
preClearMap.preClearMap();
registries.forEach((registry -> {
if (isDuplicateInMap(registry).isEmpty()) {
insertDepoBalanceRegister(registry);

View file

@ -16,12 +16,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.clearing.registry.util.PreClearMap;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.validation.IValidator;
import java.text.MessageFormat;
import java.time.Instant;
import java.util.Collection;
import java.util.Map;
@ -41,6 +40,7 @@ public class DepoPaymentInstructionRegisterService extends QueueConsumer impleme
private final Imdg<TradingClearingRegistry> tradingClearingRegistryMap;
private final Imdg<InOutDirectionDictionary> InOutDirectionDictionaryMap;
private final Function<Map<String, ?>, IValidator> fieldValuesValidator;
private final PreClearMap<DepoPaymentInstructionRegister> preClearMap;
@Autowired
public DepoPaymentInstructionRegisterService(Consumer<String, Object> kafkaQueue,
@ -54,6 +54,7 @@ public class DepoPaymentInstructionRegisterService extends QueueConsumer impleme
this.tradingClearingRegistryMap = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
this.InOutDirectionDictionaryMap = imdgProvider.getImdg(IMDGDistributedNames.Map_InOutDirectionDictionary, InOutDirectionDictionary.class);
this.fieldValuesValidator = fieldValuesValidator;
preClearMap = PreClearMap.instanceForInstantField(depoPaymentInstructionRegisterMap, "created");
}
@Override
@ -67,15 +68,19 @@ public class DepoPaymentInstructionRegisterService extends QueueConsumer impleme
public void depoPaymentInstructionRegisterNew(BaseRequest<LauncherCommandRequest> userRequest) {
log.debug("LauncherCommandRequest received from {}", createRegistry_GORR.topic());
Collection<Session> actualSessions = sessionMap.getCollectionObjectsByFieldValues(Map.of("section", "FOND"));
Map<String, String> sessionQuery = Map.of("section", "FOND");
Collection<Session> actualSessions = sessionMap.getCollectionObjectsByFieldValues(sessionQuery);
if (actualSessions.isEmpty()) {
log.debug("No Session with such conditions section = FOND");
return;
}
log.debug("Select {} Session by query {}", actualSessions.size(), sessionQuery);
preClearMap.preClearMap();
for (Session session : actualSessions) {
Map<String, ? extends Comparable<?>> fieldValues = Map.of("sessionId", session.getId());
Collection<PaymentInstruction> paymentInstructionBySessionId =
paymentInstructionMap.getCollectionObjectsByFieldValues(fieldValues);
log.debug("Select {} PaymentInstruction by query {}", paymentInstructionBySessionId.size(), fieldValues);
for (PaymentInstruction paymentInstruction : paymentInstructionBySessionId) {
if (isDuplicateInMap(paymentInstruction).isEmpty()) {
insertDepoPaymentInstructionRegister(paymentInstruction);

View file

@ -19,15 +19,14 @@ import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.registry.ExcludeLiabilitiesRegisterNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.clearing.registry.util.PreClearMap;
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.specific.SecuritySelector;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.time.TimeUtil;
import ru.spcex.platform.utils.validation.IValidator;
import java.text.MessageFormat;
import java.time.Instant;
import java.time.LocalDate;
import java.util.Collection;
@ -47,6 +46,7 @@ public class ExcludeLiabilitiesRegisterService extends QueueConsumer implements
private final Imdg<ru.clearing.classes.statics.data.company.CompanySymbols> companySymbolMap;
private final SecuritySelector<Security> securitySelector;
private final Function<Map<String, ?>, IValidator> fieldValuesValidator;
private final PreClearMap<ExcludeLiabilitiesRegister> preClearMap;
@Autowired
public ExcludeLiabilitiesRegisterService(Consumer<String, Object> kafkaQueue,
@ -60,6 +60,7 @@ public class ExcludeLiabilitiesRegisterService extends QueueConsumer implements
this.companySymbolMap = imdgProvider.getImdg(IMDGDistributedNames.Map_CompanySymbols, ru.clearing.classes.statics.data.company.CompanySymbols.class);
this.securitySelector = new SecuritySelector<>(imdgProvider, Security.class);
this.fieldValuesValidator = fieldValuesValidator;
preClearMap = PreClearMap.instanceForInstantField(excludeLiabilitiesRegisterMap, "created");
}
@Override
@ -82,6 +83,8 @@ public class ExcludeLiabilitiesRegisterService extends QueueConsumer implements
log.debug("No registry with such conditions {}", sqlForRegistries);
return;
}
log.debug("Select {} Registry by query {}", registries.size(), sqlForRegistries);
preClearMap.preClearMap();
registries.forEach((registry -> {
if (isDuplicateInMap(registry).isEmpty()) {
insertExcludeLiabilitiesRegister(registry);

View file

@ -18,14 +18,13 @@ import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.registry.ExecutionRegisterOnSaveRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.clearing.registry.util.PreClearMap;
import ru.spcex.platform.enumeration.CompanySymbol;
import ru.spcex.platform.enumeration.MoneyFlowSide;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.validation.IValidator;
import java.text.MessageFormat;
import java.time.Instant;
import java.time.LocalDate;
import java.util.Collection;
@ -46,6 +45,7 @@ public class ExecutionRegisterService extends QueueConsumer implements Initializ
private final Imdg<ExecutionDeposit> executionDepositMap;
private final Imdg<TradingClearingRegistry> tradingClearingRegistryMap;
private final Function<Map<String, ?>, IValidator> fieldValuesValidator;
private final PreClearMap<ExecutionRegister> preClearMap;
@Autowired
public ExecutionRegisterService(Consumer<String, Object> kafkaQueue,
@ -59,6 +59,7 @@ public class ExecutionRegisterService extends QueueConsumer implements Initializ
this.executionDepositMap = imdgProvider.getImdg(Map_ExecutionDeposit, ExecutionDeposit.class);
this.tradingClearingRegistryMap = imdgProvider.getImdg(Map_TradingClearingRegistry, TradingClearingRegistry.class);
this.fieldValuesValidator = fieldValuesValidator;
preClearMap = PreClearMap.instanceForLocalDateField(executionRegisterMap, "clearingDate");
}
@Override
@ -77,6 +78,9 @@ public class ExecutionRegisterService extends QueueConsumer implements Initializ
log.debug("ExecutionRegisterNewRequest received");
Collection<ExecutionFond> collectionFond = executionFondMap.getCollectionObjectsByFieldValues(Map.of("clearingDate", localDateNow));
Collection<ExecutionDeposit> collectionDeposit = executionDepositMap.getCollectionObjectsByFieldValues(Map.of("clearingDate", localDateNow));
log.debug("Found {} ExecutionFond and {} ExecutionDeposit on clearingDate={}",
collectionFond.size(), collectionDeposit.size(), localDateNow);
preClearMap.preClearMap();
collectionFond.forEach(x -> {
if (isDuplicateInMap(x).isEmpty()) {
insertExecutionRegister(true, x);

View file

@ -18,17 +18,16 @@ import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.registry.LiabilitiesRegisterNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.clearing.registry.util.PreClearMap;
import ru.spcex.platform.enumeration.CompanySymbol;
import ru.spcex.platform.enumeration.RegistryDesignation;
import ru.spcex.platform.enumeration.RegistryInstrumentType;
import ru.spcex.platform.enumeration.RegistryUnit;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.time.TimeUtil;
import ru.spcex.platform.utils.validation.IValidator;
import java.text.MessageFormat;
import java.time.Instant;
import java.util.Collection;
import java.util.Map;
@ -47,6 +46,7 @@ public class LiabilitiesRegisterService extends QueueConsumer implements Initial
private final Imdg<CompanySymbols> companySymbolsMap;
private final Imdg<MoneyMarketSecurity> moneyMarketSecurityMap;
private final Function<Map<String, ?>, IValidator> fieldValuesValidator;
private final PreClearMap<LiabilitiesRegister> preClearMap;
@Autowired
public LiabilitiesRegisterService(Consumer<String, Object> kafkaQueue,
@ -60,6 +60,7 @@ public class LiabilitiesRegisterService extends QueueConsumer implements Initial
this.companySymbolsMap = imdgProvider.getImdg(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class);
this.moneyMarketSecurityMap = imdgProvider.getImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, MoneyMarketSecurity.class);
this.fieldValuesValidator = fieldValuesValidator;
preClearMap = PreClearMap.instanceForInstantField(liabilitiesRegisterMap, "created");
}
@Override
@ -78,7 +79,8 @@ public class LiabilitiesRegisterService extends QueueConsumer implements Initial
log.debug("LiabilitiesRegisterNewRequest received from {}", Consts.REGISTRY_LIABILITIES_REGISTER_NEW);
String sqlConditionForRegistry = getSqlForRegistries();
Collection<Registry> registries = registryMap.getCollectionObjectsBySQL(sqlConditionForRegistry);
//checking for duplicates
log.debug("Found {} Registry by query {}", registries.size(), sqlConditionForRegistry);
preClearMap.preClearMap();
registries.forEach((registry -> {
if (isDuplicateInMap(registry).isEmpty()) {
insertLiabilitiesRegister(registry);

View file

@ -16,6 +16,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.clearing.registry.util.PreClearMap;
import ru.spcex.platform.enumeration.CompanySymbol;
import ru.spcex.platform.enumeration.RegistryTradingParams;
import ru.spcex.platform.imdg.api.Imdg;
@ -23,13 +24,10 @@ 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.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.validation.IValidator;
import java.text.MessageFormat;
import java.time.Instant;
import java.util.Collection;
import java.util.HashMap;
import java.util.Map;
import java.util.Optional;
import java.util.function.Function;
@ -46,6 +44,7 @@ public class MoneyBalanceRegisterService extends QueueConsumer implements Initia
private final Imdg<Company> companyMap;
private final Imdg<Account> accountMap;
private final Function<Map<String, ?>, IValidator> fieldValuesValidator;
private final PreClearMap<MoneyBalanceRegister> preClearMap;
@Autowired
public MoneyBalanceRegisterService(Consumer<String, Object> kafkaQueue,
@ -59,6 +58,7 @@ public class MoneyBalanceRegisterService extends QueueConsumer implements Initia
this.companyMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
this.accountMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class);
this.fieldValuesValidator = fieldValuesValidator;
preClearMap = PreClearMap.instanceForInstantField(moneyBalanceRegisterMap, "created");
}
@Override
@ -73,11 +73,12 @@ public class MoneyBalanceRegisterService extends QueueConsumer implements Initia
log.debug("LauncherCommandRequest received from {}", createRegistry_GBRR.topic());
ImdgPredicate sqlConditionForRegistry = getSqlForRegistries(userRequest.getRequestPayload().getCompanyId());
Collection<Registry> registries = registryMap.getCollectionObjectsByPredicate(sqlConditionForRegistry);
if (registries == null) {
if (registries.isEmpty()) {
log.debug("No registry with such conditions {}", sqlConditionForRegistry);
return;
}
//todo надо группировать!
log.debug("Select {} Registry by query {}", registries.size(), sqlConditionForRegistry);
preClearMap.preClearMap();
registries.forEach((registry -> {
if (isDuplicateInMap(registry).isEmpty()) {
insertMoneyBalanceRegister(registry);

View file

@ -14,12 +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.clearing.registry.util.PreClearMap;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.validation.IValidator;
import java.text.MessageFormat;
import java.time.Instant;
import java.util.Collection;
import java.util.Map;
@ -36,6 +35,7 @@ public class MoneyPaymentInstructionRegisterService extends QueueConsumer implem
private final Imdg<PaymentInstruction> paymentInstructionMap;
private final Imdg<Session> sessionMap;
private final Function<Map<String, ?>, IValidator> fieldValuesValidator;
private final PreClearMap<MoneyPaymentInstructionRegister> preClearMap;
@Autowired
public MoneyPaymentInstructionRegisterService(Consumer<String, Object> kafkaQueue,
@ -47,6 +47,7 @@ public class MoneyPaymentInstructionRegisterService extends QueueConsumer implem
this.paymentInstructionMap = imdgProvider.getImdg(IMDGDistributedNames.Map_PaymentInstruction, PaymentInstruction.class);
this.sessionMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class);
this.fieldValuesValidator = fieldValuesValidator;
preClearMap = PreClearMap.instanceForLocalDateField(moneyPaymentInstructionRegisterMap, "clearingDate");
}
@Override
@ -61,10 +62,13 @@ public class MoneyPaymentInstructionRegisterService extends QueueConsumer implem
log.debug("LauncherCommandRequest received from {}", TOPIC_NAME);
Map<String, ? extends Comparable<?>> fieldValuesSectionMkr = Map.of("section", "MKR");
Collection<Session> actualSessions = sessionMap.getCollectionObjectsByFieldValues(fieldValuesSectionMkr);
log.debug("Selected {} Session by query {}", actualSessions.size(), fieldValuesSectionMkr);
preClearMap.preClearMap();
actualSessions.forEach(session -> {
Map<String, ? extends Comparable<?>> fieldValuesSessionId = Map.of("sessionId", session.getId());
Collection<PaymentInstruction> paymentInstructionBySessionId =
paymentInstructionMap.getCollectionObjectsByFieldValues(fieldValuesSessionId);
log.debug("Selected {} PaymentInstruction by query {}", paymentInstructionBySessionId.size(), fieldValuesSessionId);
for (PaymentInstruction paymentInstruction : paymentInstructionBySessionId) {
if (isDuplicateInMap(paymentInstruction).isEmpty()) {
insertMoneyPaymentInstructionRegister(paymentInstruction);

View file

@ -0,0 +1,79 @@
package ru.spcex.clearing.registry.util;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import ru.clearing.classes.objects.BusinessObject;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
import ru.spcex.platform.utils.time.TimeUtil;
import java.time.Instant;
import java.time.LocalDate;
import java.time.temporal.ChronoUnit;
import java.util.Collection;
import java.util.Objects;
public class PreClearMap<T extends BusinessObject> {
protected final Logger log = LoggerFactory.getLogger(getClass());
protected final Imdg<T> map;
protected final PredicateBuilder predicateBuilder;
public PreClearMap(Imdg<T> map, PredicateBuilder predicateBuilder) {
this.map = map;
this.predicateBuilder = predicateBuilder;
}
public void preClearMap() {
ImdgPredicate query = predicateBuilder.query(map.predicateBuilder(), LocalDate.now());
Collection<T> toDelete = map.getCollectionObjectsByPredicate(query);
log.debug("Prepare {} to delete by query \"{}\" from map {}", toDelete.size(), query, map.getMapName());
for (T item : toDelete)
map.delete(item);
}
public static <T extends BusinessObject> PreClearMap<T> instanceForLocalDateField(Imdg<T> map, String fieldName) {
Objects.requireNonNull(map);
Objects.requireNonNull(fieldName);
return new PreClearMap<>(map, new PredicateBuilder() {
@Override
public ImdgPredicate query(ImdgPredicateBuilder pb, LocalDate date) {
return pb.equals(fieldName, date);
}
@Override
public String toString() {
return "LocalDate predicate on " + fieldName;
}
});
}
public static <T extends BusinessObject> PreClearMap<T> instanceForInstantField(Imdg<T> map, String fieldName) {
Objects.requireNonNull(map);
Objects.requireNonNull(fieldName);
return new PreClearMap<>(map, new PredicateBuilder() {
@Override
public ImdgPredicate query(ImdgPredicateBuilder pb, LocalDate date) {
// Instant atFrom = TimeUtil.localDateToInstant(date);
// Instant atTill = TimeUtil.localDateToInstant(date.plusDays(1));
Instant atFrom = TimeUtil.localDateToInstant(date).truncatedTo(ChronoUnit.DAYS);
Instant atTill = TimeUtil.localDateToInstant(date.plusDays(1)).plus(1, ChronoUnit.DAYS).truncatedTo(ChronoUnit.DAYS);
return pb.and(pb.greatEqual(fieldName, atFrom), pb.less(fieldName, atTill));
}
@Override
public String toString() {
return "Instant predicate on " + fieldName;
}
});
}
public interface PredicateBuilder {
ImdgPredicate query(ImdgPredicateBuilder pb, LocalDate onDate);
}
@Override
public String toString() {
return "PreClearMap{map=" + map.getMapName() + ", predicateBuilder=" + predicateBuilder + "}";
}
}

View file

@ -0,0 +1,50 @@
package ru.spcex.clearing.registry.util;
import org.junit.jupiter.api.Test;
import org.mockito.Mockito;
import ru.clearing.classes.objects.BusinessObject;
import ru.spcex.platform.imdg.api.Imdg;
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.LocalDate;
import static org.junit.jupiter.api.Assertions.*;
class PreClearMapTest {
@Test
void instanceForLocalDateField() {
LocalDate dt = LocalDate.of(2023, 9, 29);
Imdg<BusinessObject> map = Mockito.mock(Imdg.class);
Mockito.when(map.getMapName()).thenReturn("Map_Test");
ImdgPredicateBuilder pb = new ImdgPredicateBuilderHazelcast();
Mockito.when(map.predicateBuilder()).thenReturn(pb);
PreClearMap<BusinessObject> pcm = PreClearMap.instanceForLocalDateField(map, "dateField");
assertEquals("PreClearMap{map=Map_Test, predicateBuilder=LocalDate predicate on dateField}", pcm.toString());
ImdgPredicate predicate = pcm.predicateBuilder.query(pb, dt);
assertEquals("dateField=2023-09-29", predicate.toString());
}
@Test
void instanceForInstantField() {
LocalDate dt = LocalDate.of(2023, 9, 29);
Imdg<BusinessObject> map = Mockito.mock(Imdg.class);
Mockito.when(map.getMapName()).thenReturn("Map_Test");
ImdgPredicateBuilder pb = new ImdgPredicateBuilderHazelcast();
Mockito.when(map.predicateBuilder()).thenReturn(pb);
PreClearMap<BusinessObject> pcm = PreClearMap.instanceForInstantField(map, "dateField");
assertEquals("PreClearMap{map=Map_Test, predicateBuilder=Instant predicate on dateField}", pcm.toString());
ImdgPredicate predicate = pcm.predicateBuilder.query(pb, dt);
assertEquals("(dateField>=2023-09-28T00:00:00Z AND dateField<2023-09-30T00:00:00Z)", predicate.toString());
// String expectedQuery = "(dateField>=" + TimeUtil.localDateToInstant(dt) + " AND dateField<" + TimeUtil.localDateToInstant(dt.plusDays(1)) + ")";
// assertEquals("(dateField>=2023-09-28T21:00:00Z AND dateField<2023-09-29T21:00:00Z)", predicate.toString());
// assertEquals(expectedQuery, predicate.toString());
}
}