diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FinishingSession.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FinishingSession.java index e703eddfd..c96ea36f1 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FinishingSession.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FinishingSession.java @@ -4,6 +4,7 @@ import java.time.Instant; import java.time.LocalDate; import java.util.ArrayList; import java.util.Collection; +import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Objects; @@ -110,13 +111,22 @@ public class FinishingSession implements ISessionStage { protected StageResult finishingSession(Long sessionId, String pr) { Instant now = Instant.now(); //установка CLRD для обработанных регистров - Collection claimsAndLiabilities = selectClaimsAndLiabilities(sessionId); - claimsAndLiabilities.forEach(rgs -> { - rgs.setRegistryStatus(RegistryStatus.CLRD.getKey()); - rgs.setUpdated(now); - registryImdg.update(rgs); - - }); + Collection allLiabilitiesAndOKClaims = selectClaimsAndLiabilities(sessionId); + Collection claimsAndLiabilities = allLiabilitiesAndOKClaims + .stream() + .filter(rgs -> RegistryDesignation.T.equalsByKey(rgs.getRegistryDesignation()) + || RegistryStatus.OK.equalsByKey(rgs.getRegistryStatus())) + .toList(); + { + Map 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()); // Отметка статуса в Execution* @@ -124,7 +134,7 @@ public class FinishingSession implements ISessionStage { if (session == null) { log.error("Session id={} not found. Can not update Execution's.", sessionId); } else { - Collection oRegs = findORegistryBySessionId(sessionId); + Collection oRegs = findORegistryBySessionId(allLiabilitiesAndOKClaims); Map obligationStatusByGroupId; obligationStatusByGroupId = oRegs.stream() .map(rgs -> new Pair<>(rgs.getGroupId(), getEnumByKey(RegistryStatus.class, rgs.getRegistryStatus()))) @@ -140,9 +150,9 @@ public class FinishingSession implements ISessionStage { }; }; Set exchangeExecutionIdsPreviousDay = oRegs.stream() - .filter(rgs -> Objects.nonNull(rgs.getSettlementDate())) + //.filter(rgs -> Objects.nonNull(rgs.getSettlementDate())) .filter(rgs -> Objects.nonNull(rgs.getTradingDate())) - .filter(rgs -> rgs.getSettlementDate().compareTo(rgs.getTradingDate()) > 0) + .filter(rgs -> rgs.getSettlementDate().isAfter(rgs.getTradingDate())) .map(Registry::getGroupId) .collect(Collectors.toSet()); log.trace("obligations number with second leg in the past: {}", exchangeExecutionIdsPreviousDay.size()); @@ -150,6 +160,7 @@ public class FinishingSession implements ISessionStage { int notAllowed = 0; SessionType sessionType = getEnumByKey(SessionType.class, session.getSessionType()); if (IEnumKey.contains(sessionType, SessionType.FINL, SessionType.MEDM, SessionType.UNIT)) { + Map depos = new HashMap<>(); Collection executions = findExecutionDepositBySessionId(sessionId); log.trace("loaded {} ExecutionDeposits for sessionId {}", executions.size(), sessionId); for (ExecutionDeposit execution : executions) { @@ -166,8 +177,10 @@ public class FinishingSession implements ISessionStage { } execution.setCoverageStatus(toStatus.getKey()); 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 {}", sessionId, 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)) { Collection executions = findExecutionFondBySessionId(sessionId, sessionType); + Map execCurr = new HashMap<>(); + Map execFond = new HashMap<>(); log.trace("loaded {} Executions for sessionId {}", executions.size(), sessionId); for (ExecutionCommon execution : executions) { CoverageStatus toStatus = statusByExchangeId.apply(execution.getExchangeExecutionId()); @@ -193,11 +208,15 @@ public class FinishingSession implements ISessionStage { execution.setCoverageStatus(toStatus.getKey()); execution.setUpdated(Instant.now()); if (execution.type().equals(ExecutionType.ExecutionCurrency)) { - executionCurrencyImdg.update((ExecutionCurrency) execution); + execCurr.put(execution.getId(), (ExecutionCurrency) execution); + //executionCurrencyImdg.update((ExecutionCurrency) execution); } 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 {}", sessionId, executions.size(), @@ -281,14 +300,16 @@ public class FinishingSession implements ISessionStage { protected Collection selectClaimsAndLiabilities(Long sessionId) { ImdgPredicateBuilder pb = registryImdg.predicateBuilder(); ImdgPredicate prdct = pb.and( - pb.equals("sessionId", sessionId), - pb.equals("registryStatus", RegistryStatus.OK.getKey()), - pb.or( - pb.equals("registryDesignation", RegistryDesignation.O.getKey()), - pb.equals("registryDesignation", RegistryDesignation.T.getKey()), - pb.equals("registryDesignation", RegistryDesignation.C.getKey()), - pb.equals("registryDesignation", RegistryDesignation.L.getKey()) - ) + pb.equals("sessionId", sessionId), + pb.or( + pb.and( + pb.equals("registryStatus", RegistryStatus.OK.getKey()), + pb.equals("registryDesignation", RegistryDesignation.T.getKey()) + //pb.equals("registryDesignation", RegistryDesignation.C.getKey()), + //pb.equals("registryDesignation", RegistryDesignation.L.getKey()) + ), + pb.equals("registryDesignation", RegistryDesignation.O.getKey()) + ) ); log.debug("searching claims and liabilities by predicate: {}...", prdct); Collection result = registryImdg.getCollectionObjectsByPredicate(prdct); @@ -296,20 +317,14 @@ public class FinishingSession implements ISessionStage { return result; } - protected Collection findORegistryBySessionId(Long sessionId) { - ImdgPredicateBuilder pb = registryImdg.predicateBuilder(); - ImdgPredicate prdct = pb.and( - pb.equals("sessionId", sessionId), -// registryStatus - любой статус - pb.equals("registryDesignation", RegistryDesignation.O.getKey())); - Collection result = registryImdg.getCollectionObjectsByPredicate(prdct); - int size = 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); + protected Collection findORegistryBySessionId(Collection rgss) { + List result = rgss + .stream() + .filter(rgs -> RegistryDesignation.O.equalsByKey(rgs.getRegistryDesignation())) + .filter(registry -> Objects.equals(registry.getSettlementDate(), LocalDate.now())) + .filter(registry -> Objects.equals(registry.getValueDate(), registry.getSettlementDate())) + .collect(Collectors.toList()); + log.debug("{} liabilities", result.size()); return result; } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java index 13006c37b..2076a0ea3 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java @@ -30,6 +30,7 @@ import ru.spcex.clearing.session.stage.ISessionStage; import ru.spcex.clearing.session.stage.StageResult; import ru.spcex.clearing.session.stage.Task; 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.enumeration.RegistryDesignation; import ru.spcex.platform.enumeration.RegistryInstrumentType; @@ -235,22 +236,6 @@ public class InclusionObligations implements ISessionStage { } } - private static Map castMap(Map m) { - return (Map) 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) m; -// } -// case ExecutionCurrency -> { -// return (Map) m; -// } -// } - } - private void updateExecsBatchV2(List execs) { log.debug("updating {} executions", execs.size()); execs.sort(Comparator.comparing(ExecutionCommon::type)); @@ -275,9 +260,9 @@ public class InclusionObligations implements ISessionStage { log.debug("batch size {} type {}", m.size(), execType); if (!m.isEmpty()) { switch (execType) { - case ExecutionDeposit -> executionDepositImdg.putAll(castMap(m)); - case ExecutionFond -> executionFondImdg.putAll(castMap(m)); - case ExecutionCurrency -> executionCurrImdg.putAll(castMap(m)); + case ExecutionDeposit -> executionDepositImdg.putAll(ClearingUtil.castMap(m)); + case ExecutionFond -> executionFondImdg.putAll(ClearingUtil.castMap(m)); + case ExecutionCurrency -> executionCurrImdg.putAll(ClearingUtil.castMap(m)); } } log.debug("batch size {} type {} done", m.size(), execType); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/util/ClearingUtil.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/util/ClearingUtil.java new file mode 100644 index 000000000..35294a2d2 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/util/ClearingUtil.java @@ -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 Map castMap(Map m) { + return (Map) 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) m; +// } +// case ExecutionCurrency -> { +// return (Map) m; +// } +// } + } +}