session step 10 optimization - C/L

This commit is contained in:
ialbert 2024-08-02 18:54:21 +03:00
parent b74ebc3b0e
commit d7b3189a80
3 changed files with 76 additions and 54 deletions

View file

@ -4,6 +4,7 @@ import java.time.Instant;
import java.time.LocalDate; import java.time.LocalDate;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Collection; import java.util.Collection;
import java.util.HashMap;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Objects; import java.util.Objects;
@ -110,13 +111,22 @@ public class FinishingSession implements ISessionStage {
protected StageResult<?> finishingSession(Long sessionId, String pr) { protected StageResult<?> finishingSession(Long sessionId, String pr) {
Instant now = Instant.now(); Instant now = Instant.now();
//установка CLRD для обработанных регистров //установка CLRD для обработанных регистров
Collection<Registry> claimsAndLiabilities = selectClaimsAndLiabilities(sessionId); Collection<Registry> allLiabilitiesAndOKClaims = selectClaimsAndLiabilities(sessionId);
claimsAndLiabilities.forEach(rgs -> { Collection<Registry> claimsAndLiabilities = allLiabilitiesAndOKClaims
rgs.setRegistryStatus(RegistryStatus.CLRD.getKey()); .stream()
rgs.setUpdated(now); .filter(rgs -> RegistryDesignation.T.equalsByKey(rgs.getRegistryDesignation())
registryImdg.update(rgs); || RegistryStatus.OK.equalsByKey(rgs.getRegistryStatus()))
.toList();
}); {
Map<Long, Registry> rgsToUpdate = new HashMap<>();
claimsAndLiabilities.forEach(rgs -> {
rgs.setRegistryStatus(RegistryStatus.CLRD.getKey());
rgs.setUpdated(now);
// registryImdg.update(rgs);
rgsToUpdate.put(rgs.getId(), rgs);
});
registryImdg.putAll(rgsToUpdate, 200);
}
log.debug("For sessionId={} was updated {} registers for status={}", sessionId, claimsAndLiabilities.size(), RegistryStatus.CLRD.getKey()); log.debug("For sessionId={} was updated {} registers for status={}", sessionId, claimsAndLiabilities.size(), RegistryStatus.CLRD.getKey());
// Отметка статуса в Execution* // Отметка статуса в Execution*
@ -124,7 +134,7 @@ public class FinishingSession implements ISessionStage {
if (session == null) { if (session == null) {
log.error("Session id={} not found. Can not update Execution's.", sessionId); log.error("Session id={} not found. Can not update Execution's.", sessionId);
} else { } else {
Collection<Registry> oRegs = findORegistryBySessionId(sessionId); Collection<Registry> oRegs = findORegistryBySessionId(allLiabilitiesAndOKClaims);
Map<Long, RegistryStatus> obligationStatusByGroupId; Map<Long, RegistryStatus> obligationStatusByGroupId;
obligationStatusByGroupId = oRegs.stream() obligationStatusByGroupId = oRegs.stream()
.map(rgs -> new Pair<>(rgs.getGroupId(), getEnumByKey(RegistryStatus.class, rgs.getRegistryStatus()))) .map(rgs -> new Pair<>(rgs.getGroupId(), getEnumByKey(RegistryStatus.class, rgs.getRegistryStatus())))
@ -140,9 +150,9 @@ public class FinishingSession implements ISessionStage {
}; };
}; };
Set<Long> exchangeExecutionIdsPreviousDay = oRegs.stream() Set<Long> exchangeExecutionIdsPreviousDay = oRegs.stream()
.filter(rgs -> Objects.nonNull(rgs.getSettlementDate())) //.filter(rgs -> Objects.nonNull(rgs.getSettlementDate()))
.filter(rgs -> Objects.nonNull(rgs.getTradingDate())) .filter(rgs -> Objects.nonNull(rgs.getTradingDate()))
.filter(rgs -> rgs.getSettlementDate().compareTo(rgs.getTradingDate()) > 0) .filter(rgs -> rgs.getSettlementDate().isAfter(rgs.getTradingDate()))
.map(Registry::getGroupId) .map(Registry::getGroupId)
.collect(Collectors.toSet()); .collect(Collectors.toSet());
log.trace("obligations number with second leg in the past: {}", exchangeExecutionIdsPreviousDay.size()); log.trace("obligations number with second leg in the past: {}", exchangeExecutionIdsPreviousDay.size());
@ -150,6 +160,7 @@ public class FinishingSession implements ISessionStage {
int notAllowed = 0; int notAllowed = 0;
SessionType sessionType = getEnumByKey(SessionType.class, session.getSessionType()); SessionType sessionType = getEnumByKey(SessionType.class, session.getSessionType());
if (IEnumKey.contains(sessionType, SessionType.FINL, SessionType.MEDM, SessionType.UNIT)) { if (IEnumKey.contains(sessionType, SessionType.FINL, SessionType.MEDM, SessionType.UNIT)) {
Map<Long, ExecutionDeposit> depos = new HashMap<>();
Collection<ExecutionDeposit> executions = findExecutionDepositBySessionId(sessionId); Collection<ExecutionDeposit> executions = findExecutionDepositBySessionId(sessionId);
log.trace("loaded {} ExecutionDeposits for sessionId {}", executions.size(), sessionId); log.trace("loaded {} ExecutionDeposits for sessionId {}", executions.size(), sessionId);
for (ExecutionDeposit execution : executions) { for (ExecutionDeposit execution : executions) {
@ -166,8 +177,10 @@ public class FinishingSession implements ISessionStage {
} }
execution.setCoverageStatus(toStatus.getKey()); execution.setCoverageStatus(toStatus.getKey());
execution.setUpdated(Instant.now()); execution.setUpdated(Instant.now());
executionDepositImdg.update(execution); depos.put(execution.getId(), execution);
// executionDepositImdg.update(execution);
} }
executionDepositImdg.putAll(depos, 200);
log.trace("By sessionId={} processed {} ExecutionDeposit: {} allowed, {} denied, skipped {}", log.trace("By sessionId={} processed {} ExecutionDeposit: {} allowed, {} denied, skipped {}",
sessionId, sessionId,
executions.size(), executions.size(),
@ -177,6 +190,8 @@ public class FinishingSession implements ISessionStage {
} }
if (IEnumKey.contains(sessionType, SessionType.CURR, SessionType.TRDT, SessionType.IPOB, SessionType.IPO0, SessionType.IPOT, SessionType.UNIT)) { if (IEnumKey.contains(sessionType, SessionType.CURR, SessionType.TRDT, SessionType.IPOB, SessionType.IPO0, SessionType.IPOT, SessionType.UNIT)) {
Collection<? extends ExecutionCommon> executions = findExecutionFondBySessionId(sessionId, sessionType); Collection<? extends ExecutionCommon> executions = findExecutionFondBySessionId(sessionId, sessionType);
Map<Long, ExecutionCurrency> execCurr = new HashMap<>();
Map<Long, ExecutionFond> execFond = new HashMap<>();
log.trace("loaded {} Executions for sessionId {}", executions.size(), sessionId); log.trace("loaded {} Executions for sessionId {}", executions.size(), sessionId);
for (ExecutionCommon execution : executions) { for (ExecutionCommon execution : executions) {
CoverageStatus toStatus = statusByExchangeId.apply(execution.getExchangeExecutionId()); CoverageStatus toStatus = statusByExchangeId.apply(execution.getExchangeExecutionId());
@ -193,11 +208,15 @@ public class FinishingSession implements ISessionStage {
execution.setCoverageStatus(toStatus.getKey()); execution.setCoverageStatus(toStatus.getKey());
execution.setUpdated(Instant.now()); execution.setUpdated(Instant.now());
if (execution.type().equals(ExecutionType.ExecutionCurrency)) { if (execution.type().equals(ExecutionType.ExecutionCurrency)) {
executionCurrencyImdg.update((ExecutionCurrency) execution); execCurr.put(execution.getId(), (ExecutionCurrency) execution);
//executionCurrencyImdg.update((ExecutionCurrency) execution);
} else { } else {
executionFondImdg.update((ExecutionFond) execution); execFond.put(execution.getId(), (ExecutionFond) execution);
//executionFondImdg.update((ExecutionFond) execution);
} }
} }
executionCurrencyImdg.putAll(execCurr, 200);
executionFondImdg.putAll(execFond, 200);
log.trace("By sessionId={} processed {} ExecutionFond: allowed {}, denied {}, skipped {}", log.trace("By sessionId={} processed {} ExecutionFond: allowed {}, denied {}, skipped {}",
sessionId, sessionId,
executions.size(), executions.size(),
@ -281,14 +300,16 @@ public class FinishingSession implements ISessionStage {
protected Collection<Registry> selectClaimsAndLiabilities(Long sessionId) { protected Collection<Registry> selectClaimsAndLiabilities(Long sessionId) {
ImdgPredicateBuilder pb = registryImdg.predicateBuilder(); ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
ImdgPredicate prdct = pb.and( ImdgPredicate prdct = pb.and(
pb.equals("sessionId", sessionId), pb.equals("sessionId", sessionId),
pb.equals("registryStatus", RegistryStatus.OK.getKey()), pb.or(
pb.or( pb.and(
pb.equals("registryDesignation", RegistryDesignation.O.getKey()), pb.equals("registryStatus", RegistryStatus.OK.getKey()),
pb.equals("registryDesignation", RegistryDesignation.T.getKey()), pb.equals("registryDesignation", RegistryDesignation.T.getKey())
pb.equals("registryDesignation", RegistryDesignation.C.getKey()), //pb.equals("registryDesignation", RegistryDesignation.C.getKey()),
pb.equals("registryDesignation", RegistryDesignation.L.getKey()) //pb.equals("registryDesignation", RegistryDesignation.L.getKey())
) ),
pb.equals("registryDesignation", RegistryDesignation.O.getKey())
)
); );
log.debug("searching claims and liabilities by predicate: {}...", prdct); log.debug("searching claims and liabilities by predicate: {}...", prdct);
Collection<Registry> result = registryImdg.getCollectionObjectsByPredicate(prdct); Collection<Registry> result = registryImdg.getCollectionObjectsByPredicate(prdct);
@ -296,20 +317,14 @@ public class FinishingSession implements ISessionStage {
return result; return result;
} }
protected Collection<Registry> findORegistryBySessionId(Long sessionId) { protected Collection<Registry> findORegistryBySessionId(Collection<Registry> rgss) {
ImdgPredicateBuilder pb = registryImdg.predicateBuilder(); List<Registry> result = rgss
ImdgPredicate prdct = pb.and( .stream()
pb.equals("sessionId", sessionId), .filter(rgs -> RegistryDesignation.O.equalsByKey(rgs.getRegistryDesignation()))
// registryStatus - любой статус .filter(registry -> Objects.equals(registry.getSettlementDate(), LocalDate.now()))
pb.equals("registryDesignation", RegistryDesignation.O.getKey())); .filter(registry -> Objects.equals(registry.getValueDate(), registry.getSettlementDate()))
Collection<Registry> result = registryImdg.getCollectionObjectsByPredicate(prdct); .collect(Collectors.toList());
int size = result.size(); log.debug("{} liabilities", result.size());
result = result.stream()
.filter(registry -> Objects.equals(registry.getSettlementDate(), LocalDate.now()))
.filter(registry -> Objects.equals(registry.getValueDate(), registry.getSettlementDate()))
.collect(Collectors.toList());
log.debug("found {} liabilities by query {} and {} liabilities after filtering ValueDate==SettlementDate",
size, result.size(), prdct);
return result; return result;
} }

View file

@ -30,6 +30,7 @@ import ru.spcex.clearing.session.stage.ISessionStage;
import ru.spcex.clearing.session.stage.StageResult; import ru.spcex.clearing.session.stage.StageResult;
import ru.spcex.clearing.session.stage.Task; import ru.spcex.clearing.session.stage.Task;
import ru.spcex.clearing.session.stage.task.InclusionToPoolPayload; import ru.spcex.clearing.session.stage.task.InclusionToPoolPayload;
import ru.spcex.clearing.util.ClearingUtil;
import ru.spcex.platform.classes.base.interfaces.ExecutionType; import ru.spcex.platform.classes.base.interfaces.ExecutionType;
import ru.spcex.platform.enumeration.RegistryDesignation; import ru.spcex.platform.enumeration.RegistryDesignation;
import ru.spcex.platform.enumeration.RegistryInstrumentType; import ru.spcex.platform.enumeration.RegistryInstrumentType;
@ -235,22 +236,6 @@ public class InclusionObligations implements ISessionStage {
} }
} }
private static <T extends ExecutionCommon> Map<Long, T> castMap(Map<Long, ExecutionCommon> m) {
return (Map<Long, T>) m;
// ExecutionType type = m.values().stream().map(e -> e.type()).findFirst().orElse(null);
// if (type == null) {
// return Collections.emptyMap();
// }
// switch (type) {
// case ExecutionDeposit -> {
// return (Map<Long, T>) m;
// }
// case ExecutionCurrency -> {
// return (Map<Long, T>) m;
// }
// }
}
private void updateExecsBatchV2(List<ExecutionCommon> execs) { private void updateExecsBatchV2(List<ExecutionCommon> execs) {
log.debug("updating {} executions", execs.size()); log.debug("updating {} executions", execs.size());
execs.sort(Comparator.comparing(ExecutionCommon::type)); execs.sort(Comparator.comparing(ExecutionCommon::type));
@ -275,9 +260,9 @@ public class InclusionObligations implements ISessionStage {
log.debug("batch size {} type {}", m.size(), execType); log.debug("batch size {} type {}", m.size(), execType);
if (!m.isEmpty()) { if (!m.isEmpty()) {
switch (execType) { switch (execType) {
case ExecutionDeposit -> executionDepositImdg.putAll(castMap(m)); case ExecutionDeposit -> executionDepositImdg.putAll(ClearingUtil.castMap(m));
case ExecutionFond -> executionFondImdg.putAll(castMap(m)); case ExecutionFond -> executionFondImdg.putAll(ClearingUtil.castMap(m));
case ExecutionCurrency -> executionCurrImdg.putAll(castMap(m)); case ExecutionCurrency -> executionCurrImdg.putAll(ClearingUtil.castMap(m));
} }
} }
log.debug("batch size {} type {} done", m.size(), execType); log.debug("batch size {} type {} done", m.size(), execType);

View file

@ -0,0 +1,22 @@
package ru.spcex.clearing.util;
import java.util.Map;
import ru.clearing.classes.statics.data.execution.ExecutionCommon;
public class ClearingUtil {
public static <T extends ExecutionCommon> Map<Long, T> castMap(Map<Long, ExecutionCommon> m) {
return (Map<Long, T>) m;
// ExecutionType type = m.values().stream().map(e -> e.type()).findFirst().orElse(null);
// if (type == null) {
// return Collections.emptyMap();
// }
// switch (type) {
// case ExecutionDeposit -> {
// return (Map<Long, T>) m;
// }
// case ExecutionCurrency -> {
// return (Map<Long, T>) m;
// }
// }
}
}