diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/component/predicate/cash/registry/AssetByObligationCashedPredicate.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/component/predicate/cash/registry/AssetByObligationCashedPredicate.java new file mode 100644 index 000000000..ab7870b6b --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/component/predicate/cash/registry/AssetByObligationCashedPredicate.java @@ -0,0 +1,50 @@ +package ru.spcex.clearing.component.predicate.cash.registry; + +import ru.clearing.classes.statics.data.registry.Registry; +import ru.spcex.clearing.service.cash.impl.RegistryCashA__F; +import ru.spcex.platform.enumeration.RegistryInstrumentType; +import ru.spcex.platform.enumeration.RegistryTradingParams; +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.IEnumKey; + +public class AssetByObligationCashedPredicate { + private final RegistryCashA__F.CustomKey key; + private final RegistryTradingParams params; + + public static AssetByObligationCashedPredicate getPredicate(Registry rgs) { + return new AssetByObligationCashedPredicate(rgs); + } + + public ImdgPredicate cashed(ImdgPredicateBuilder pb, RegistryCashA__F cash) { + ImdgPredicate prdct = pb.and( + pb.sql(RegistryCodeSqlBuilder.getInstance(params).build()), + pb.equals("tradingClearingRegistryId", key.tcrId()), + pb.equals("securitySymbol", key.secSymbol()), + pb.equals("companyId", key.cmpId()) + ); + if (cash != null) { + return pb.cashed(prdct, cash, key); + } else { + return prdct; + } + } + + private AssetByObligationCashedPredicate(Registry rgs) { + if (rgs == null) throw new IllegalStateException("cannot obtain cash for null obligation"); + RegistryTradingParams params; + if (IEnumKey.getEnumByKey(RegistryInstrumentType.class, rgs.getRegistryInstrumentType()) == RegistryInstrumentType.S) { + params = RegistryTradingParams.AS_F; + } else if (IEnumKey.getEnumByKey(RegistryInstrumentType.class, rgs.getRegistryInstrumentType()) == RegistryInstrumentType.M) { + params = RegistryTradingParams.AM_F; + } else { + throw new IllegalStateException("unknown RegistryInstrumentType '%s' for rgs.id=%d".formatted( + rgs.getRegistryInstrumentType(), + rgs.getId() + )); + } + this.params = params; + this.key = RegistryCashA__F.key(rgs); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/component/predicate/cash/registry/AssetByTcrCompanyAccountCashedPredicate.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/component/predicate/cash/registry/AssetByTcrCompanyAccountCashedPredicate.java new file mode 100644 index 000000000..426bf4d0a --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/component/predicate/cash/registry/AssetByTcrCompanyAccountCashedPredicate.java @@ -0,0 +1,37 @@ +package ru.spcex.clearing.component.predicate.cash.registry; + +import ru.clearing.classes.statics.data.registry.Registry; +import ru.spcex.clearing.service.cash.impl.RegistryCashAssetCash; +import ru.spcex.platform.enumeration.RegistryTradingParams; +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; + +public class AssetByTcrCompanyAccountCashedPredicate { + private final RegistryCashAssetCash.CustomKey key; + private final RegistryTradingParams params; + + public static AssetByTcrCompanyAccountCashedPredicate getPredicate(Registry rgs, RegistryTradingParams params) { + return new AssetByTcrCompanyAccountCashedPredicate(rgs, params); + } + + public ImdgPredicate cashed(ImdgPredicateBuilder pb, RegistryCashAssetCash cash) { + ImdgPredicate prdct = pb.and( + pb.sql(RegistryCodeSqlBuilder.getInstance(params).build()), + pb.equals("tradingClearingRegistryId", key.tcrId()), + pb.equals("companyId", key.cmpId()), + pb.equals("accountId", key.accId()), + pb.equals("securityId", key.secId()) + ); + if (cash != null) { + return pb.cashed(prdct, cash, key); + } else { + return prdct; + } + } + + private AssetByTcrCompanyAccountCashedPredicate(Registry rgs, RegistryTradingParams params) { + this.params = params; + this.key = RegistryCashAssetCash.key(rgs, params); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/cash/impl/RegistryCashA__F.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/cash/impl/RegistryCashA__F.java index ab4bfeb0d..751a78213 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/cash/impl/RegistryCashA__F.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/cash/impl/RegistryCashA__F.java @@ -15,7 +15,7 @@ public class RegistryCashA__F extends CashV2 registryImdg; + private final ImdgId idGen; private final RegistryManager registryManager; private final IMessageResolver msgResolver = new SimpleMessageResolver(); private SessionType sessionType; @@ -75,7 +81,8 @@ public class InspectionObligationsV2 implements ISessionStage { RegistryManager registryManager, AssetTBFProcessing assets, GatewayRequester gateway, PlanBalanceCalc planBalanceCalc) { - this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); + this.registryImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_Registry, Registry.class, null); + this.idGen = imdgProvider.getImdgIdGenerator(); this.registryManager = registryManager; this.assets = assets; this.gateway = gateway; @@ -121,6 +128,7 @@ public class InspectionObligationsV2 implements ISessionStage { RegistryStatus.POOL.getKey()); Collection registriesToProcess = registryImdg.getCollectionObjectsBySQL(sqlCondition); + Map rgssToStore = new HashMap<>(registriesToProcess.size()); //вторые ноги у сделок - нужна ли оптимизация? List>> registriesByGroupSorted = registriesToProcess.stream() .collect(Collectors.groupingBy(Registry::getGroupId)) .entrySet() @@ -140,8 +148,8 @@ public class InspectionObligationsV2 implements ISessionStage { }) .toList(); - List assetsToUpdate = new ArrayList<>(); log.info("found {} ({} groups) registries by sql: {}", registriesToProcess.size(), registriesByGroupSorted.size(), sqlCondition); + ImdgPredicateBuilder rgsPb = registryImdg.predicateBuilder(); GROUP: for (Map.Entry> entry : registriesByGroupSorted) { List group = entry.getValue(); @@ -159,9 +167,10 @@ public class InspectionObligationsV2 implements ISessionStage { checkResults.add(new CheckResult(obligation, false)); continue; } - String sqlAssetRegistryCondition = searchAssetsByObligationSql(obligation); - log.debug("search asset by registry.id={} sql {}", obligation.getId(), sqlAssetRegistryCondition); - Registry asset = a__fCash.getOrFind(obligation, () -> registryImdg.getFirstObjectBySQL(sqlAssetRegistryCondition)); + ImdgPredicate assetPrdct = AssetByObligationCashedPredicate + .getPredicate(obligation) + .cashed(rgsPb, a__fCash); + Registry asset = registryImdg.getFirstObjectByPredicate(assetPrdct); if (asset == null) { log.warn("Not found asset by registry.id={}", obligation.getId()); isUncovered = true; @@ -179,7 +188,7 @@ public class InspectionObligationsV2 implements ISessionStage { Runnable failGroup = () -> Stream.concat(group.stream(), getRefundDateRgsIfPresent(group).stream()) .forEach(registry -> { registry.setComment(msgResolver.resolve(RgsError.A__tNotFound)); - updateRegistryStatus(registry, RegistryStatus.FAIL); + updateRegistryStatus(registry, RegistryStatus.FAIL, rgssToStore); }); RegistryInstrumentType mOrS = IEnumKey.getEnumByKey(RegistryInstrumentType.class, rgs.getRegistryInstrumentType()); if (mOrS == null) { @@ -191,9 +200,11 @@ public class InspectionObligationsV2 implements ISessionStage { if (Section.MKR.equalsByKey(rgs.getSection()) && mOrS.equals(RegistryInstrumentType.S)) { continue; } - Optional a__t = a___Cash.getOrFindO(rgs, A__T.type(mOrS), () -> - assets.searchByTcrCompanyAccount(rgs, A__T.type(mOrS))); - if (a__t.isEmpty()) { + ImdgPredicate assetPredicate = AssetByTcrCompanyAccountCashedPredicate + .getPredicate(rgs, A__T.type(mOrS)) + .cashed(rgsPb, a___Cash); + Registry a__t = registryImdg.getFirstObjectByPredicate(assetPredicate); + if (a__t == null) { log.debug("groupId {}, {}.id={} - A**T not found. settings FAIL to group", entry.getKey(), rgs.getRegistryCode(), rgs.getId()); failGroup.run(); @@ -215,7 +226,7 @@ public class InspectionObligationsV2 implements ISessionStage { .ifPresent(chk -> chk.isUncovered = true); } } - defineStatusAndUpdateRegistry(checkResults, group); + defineStatusAndUpdateRegistry(checkResults, group, rgssToStore); if (checkResults.stream().noneMatch(checkResult -> checkResult.isUncovered)) { Instant now = Instant.now(); //по каждому регистру OM*T, TM*T, OS*T, TS*T из одной группы @@ -223,86 +234,70 @@ public class InspectionObligationsV2 implements ISessionStage { { RegistryInstrumentType mOrS = IEnumKey.getEnumByKey(RegistryInstrumentType.class, registry.getRegistryInstrumentType()); if (mOrS == null) continue; - Optional a__t = a___Cash.getOrFindO(registry, A__T.type(mOrS), () - -> assets.searchByTcrCompanyAccount(registry, A__T.type(mOrS))); - Optional a__b = a__t.flatMap(a__tFound -> a___Cash.getOrFind( - a__tFound, - A__B.type(mOrS), - () -> assets - .searchByTcrCompanyAccount(a__tFound, A__B.type(mOrS)).orElseGet(() -> copyB(a__tFound)))); -// Optional a__b = a__t.map(a__tFound -> assets -// .searchByTcrCompanyAccount(a__tFound, A__B.type(mOrS)) -// .orElseGet(() -> copyB(a__tFound))); - Optional a__f = a__t.flatMap(a__tFound -> a___Cash.getOrFind( - a__tFound, - A__F.type(mOrS), - () -> assets - .searchByTcrCompanyAccount(a__tFound, A__F.type(mOrS)).orElseGet(() -> copyF(a__tFound)))); -// Optional a__f = a__t.map(a__tFound -> assets -// .searchByTcrCompanyAccount(a__tFound, A__F.type(mOrS)) -// .orElseGet(() -> copyF(a__tFound))); - log.debug("registry {}.id={}: {}", registry.getRegistryCode(), registry.getId(), - a__t.map(r -> "A**T.id=%d/A**B.id=%d/A**F.id=%d" - .formatted(a__t.get().getId(), a__b.get().getId(), a__f.get().getId())) - .orElse("A**T/A**B/A**F not found")); - if (a__t.isPresent()) { + ImdgPredicate a__tPrdct = AssetByTcrCompanyAccountCashedPredicate + .getPredicate(registry, A__T.type(mOrS)) + .cashed(rgsPb, a___Cash); + Registry a__t = registryImdg.getFirstObjectByPredicate(a__tPrdct); + Registry a__b = null; + Registry a__f = null; + if (a__t != null) { + ImdgPredicate a__bPrdct = AssetByTcrCompanyAccountCashedPredicate + .getPredicate(registry, A__B.type(mOrS)) + .cashed(rgsPb, a___Cash); + ImdgPredicate a__fPrdct = AssetByTcrCompanyAccountCashedPredicate + .getPredicate(registry, A__F.type(mOrS)) + .cashed(rgsPb, a___Cash); + + a__b = Optional.ofNullable( + registryImdg.getFirstObjectByPredicate(a__bPrdct) + ).orElseGet(() -> copyB(a__t)); + + a__f = Optional.ofNullable( + registryImdg.getFirstObjectByPredicate(a__fPrdct) + ).orElseGet(() -> copyF(a__t)); + + } + { + String msg; + if (a__t != null) { + msg = "A**T.id=%d/A**B.id=%d/A**F.id=%d".formatted( + a__t.getId(), a__b.getId(), a__f.getId() + ); + } else { + msg = "A**T/A**B/A**F not found"; + } + log.debug("registry {}.id={}: {}", registry.getRegistryCode(), registry.getId(), msg); + } + + if (a__t != null) { BigDecimal amount = safeBD(registry.getBalance()); if (IEnumKey.getEnumByKey(RegistryDesignation.class, registry.getRegistryDesignation()) == RegistryDesignation.O) { amount = amount.negate(); } else if (IEnumKey.getEnumByKey(RegistryDesignation.class, registry.getRegistryDesignation()) == RegistryDesignation.T) { amount = amount; } - a__b.get().setSessionId(sessionId); - assetsCashing.process(a__b.get(), a__t.get(), a__f.get(), amount); - assetsToUpdate.add(a__b.get()); - assetsToUpdate.add(a__f.get()); + a__b.setSessionId(sessionId); + assetsCashing.process(a__b, a__t, a__f, amount); + rgssToStore.put(a__b.getId(), a__b); + rgssToStore.put(a__f.getId(), a__f); } } } } } - for (Registry registry : assetsToUpdate) { - if (registry.getId() != null) { - registryImdg.update(registry); - } else { - registryImdg.insert(registry); - } - } + registryImdg.putAll(rgssToStore); planBalanceCalc.stageRevision2(sessionId); return new StageResult(null, true); } - private String searchAssetsByObligationSql(Registry obligation) { - RegistryTradingParams counterRegistryTradingParams = null; - if (IEnumKey.getEnumByKey(RegistryInstrumentType.class, obligation.getRegistryInstrumentType()) == RegistryInstrumentType.S) { - counterRegistryTradingParams = new RegistryTradingParams(RegistryDesignation.A, - RegistryInstrumentType.S, - null, - RegistryUnit.F); - } else if (IEnumKey.getEnumByKey(RegistryInstrumentType.class, obligation.getRegistryInstrumentType()) == RegistryInstrumentType.M) { - counterRegistryTradingParams = new RegistryTradingParams(RegistryDesignation.A, - RegistryInstrumentType.M, - null, - RegistryUnit.F); - } - - return String.format("%s and " + - "tradingClearingRegistryId = '%s' and " + - "securitySymbol = '%s' and " + - "companyId = %d", - RegistryCodeSqlBuilder.getInstance(counterRegistryTradingParams).build(), - obligation.getTradingClearingRegistryId(), - obligation.getSecuritySymbol(), - obligation.getCompanyId()); - } - private Registry copyF(Registry rgs) { Registry rgsF = rgs.clone(); RegistryManager.zeroState(rgsF); rgsF.setRegistryUnit(RegistryUnit.F.getKey()); rgsF.setRegistryCode(RegistryUtil.clearingCode(rgsF)); + rgsF.setId(idGen.nextId()); registryImdg.insert(rgsF); log.debug("created {}.id={} by {}.id={}", rgsF.getRegistryCode(), rgsF.getId(), rgs.getRegistryCode(), rgs.getId()); return rgsF; @@ -313,12 +308,13 @@ public class InspectionObligationsV2 implements ISessionStage { RegistryManager.zeroState(rgsB); rgsB.setRegistryUnit(RegistryUnit.B.getKey()); rgsB.setRegistryCode(RegistryUtil.clearingCode(rgsB)); + rgsB.setId(idGen.nextId()); registryImdg.insert(rgsB); log.debug("created {}.id={} by {}.id={}", rgsB.getRegistryCode(), rgsB.getId(), rgs.getRegistryCode(), rgs.getId()); return rgsB; } - private void defineStatusAndUpdateRegistry(List checkResults, List registries) { + private void defineStatusAndUpdateRegistry(List checkResults, List registries, Map rgssToStore) { boolean isOneUncovered = checkResults.stream().anyMatch(checkResult -> checkResult.isUncovered); log.debug("groupId {}, uncovered: {}", registries.stream().findFirst().map(Registry::getGroupId).orElse(null), isOneUncovered); String commentErr = msgResolver.resolve(ClearingError.InsecurityObligation, checkResults.stream() @@ -334,17 +330,17 @@ public class InspectionObligationsV2 implements ISessionStage { ).findFirst(); checkResult.registry.setComment(commentErr); if (checkResult.isUncovered) { - updateRegistryStatus(checkResult.registry, uncvStatus()); - tRegistryWithSameCompany.ifPresent(registry -> updateRegistryStatus(registry, failStatus())); + updateRegistryStatus(checkResult.registry, uncvStatus(), rgssToStore); + tRegistryWithSameCompany.ifPresent(registry -> updateRegistryStatus(registry, failStatus(), rgssToStore)); } else { - updateRegistryStatus(checkResult.registry, failStatus()); - tRegistryWithSameCompany.ifPresent(registry -> updateRegistryStatus(registry, failStatus())); + updateRegistryStatus(checkResult.registry, failStatus(), rgssToStore); + tRegistryWithSameCompany.ifPresent(registry -> updateRegistryStatus(registry, failStatus(), rgssToStore)); } } //проверяем есть ли второй день для сделки (он не входит в пул, поэтому ищем отдельно) - getRefundDateRgsIfPresent(registries).forEach(rgs -> updateRegistryStatus(rgs, RegistryStatus.FAIL)); + getRefundDateRgsIfPresent(registries).forEach(rgs -> updateRegistryStatus(rgs, RegistryStatus.FAIL, rgssToStore)); } else { - registries.forEach(registry -> updateRegistryStatus(registry, RegistryStatus.OK)); + registries.forEach(registry -> updateRegistryStatus(registry, RegistryStatus.OK, rgssToStore)); } } @@ -375,6 +371,14 @@ public class InspectionObligationsV2 implements ISessionStage { registryImdg.update(registry); } + private void updateRegistryStatus(Registry registry, RegistryStatus registryStatus, Map rgss) { + log.debug("Update registry.id: {} to {}", registry.getId(), registryStatus.getKey()); + registry.setRegistryStatus(registryStatus.getKey()); + registry.setUpdated(Instant.now()); + rgss.put(registry.getId(), registry); + //registryImdg.update(registry); + } + private static final class CheckResult { private final Registry registry; private boolean isUncovered; diff --git a/platform-parent/platform-imdg-api/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/adapter/CashV2.java b/platform-parent/platform-imdg-api/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/adapter/CashV2.java index f446f0825..9dc8f23b9 100644 --- a/platform-parent/platform-imdg-api/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/adapter/CashV2.java +++ b/platform-parent/platform-imdg-api/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/adapter/CashV2.java @@ -13,10 +13,11 @@ public abstract class CashV2 { private long missRate = 0; private final Map cash; private final String name; + protected int threshold = 200_000; public CashV2() { this.cash = new HashMap<>(); - this.name = null; + this.name = "[unknown]"; } public CashV2(String name) { @@ -26,7 +27,7 @@ public abstract class CashV2 { public void clear() { log.info("cash {}: size {}, hit rate {}, miss rate {}", - name != null ? name : "[unknown]", + name, cash.size(), hitRate, missRate); @@ -38,7 +39,7 @@ public abstract class CashV2 { public void store(T obj) { if (obj == null || !needToStore(obj)) return; K key = extractKey(obj); - cash.put(key, obj); + putIfThresholdNotReached(key, obj); } protected Optional get(K key) { @@ -53,7 +54,7 @@ public abstract class CashV2 { if (by.isEmpty()) { T t = supplier.get(); if (t != null && needToStore(t)) { - cash.put(key, t); + putIfThresholdNotReached(key, t); by = Optional.of(t); } } @@ -66,7 +67,7 @@ public abstract class CashV2 { Optional foundAnew = supplier.get(); if (foundAnew.isPresent()) { if (needToStore(foundAnew.get())) { - cash.put(key, foundAnew.get()); + putIfThresholdNotReached(key, foundAnew.get()); } fromCash = foundAnew; } @@ -74,6 +75,14 @@ public abstract class CashV2 { return fromCash; } + private void putIfThresholdNotReached(K key, T value) { + if (cash.size() < threshold) { + cash.put(key, value); + } else { + log.warn("cash {}: size {} threshold reached. will not put new value in cash", name, value); + } + } + protected abstract K extractKey(T obj); protected boolean needToStore(T obj) {