remove stage 8

This commit is contained in:
etreshenkov 2023-05-25 20:15:52 +03:00
parent d3a8aaff7a
commit 61d7570b35
4 changed files with 274 additions and 24 deletions

View file

@ -2,9 +2,7 @@ package ru.spcex.clearing.config;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import ru.spcex.clearing.service.executors.AbstractExecutor;
import ru.spcex.clearing.service.executors.Sdf01Executor;
import ru.spcex.clearing.service.executors.Sdf57Executor;
import ru.spcex.clearing.service.executors.*;
import ru.spcex.platform.enumeration.SdfTable;
import java.util.HashMap;
@ -15,10 +13,14 @@ public class SdfExecutorsConfig {
@Bean("sdfExecutors")
public Map<SdfTable, AbstractExecutor<?>> executorsMap(Sdf01Executor sdf01Executor,
Sdf57Executor sdf57Executor) {
Sdf57Executor sdf57Executor,
Sdf04Executor sdf04Executor,
Sdf13Executor sdf13Executor) {
Map<SdfTable, AbstractExecutor<?>> executors = new HashMap<>();
executors.put(SdfTable.SDF_01, sdf01Executor);
executors.put(SdfTable.SDF_57, sdf57Executor);
executors.put(SdfTable.SDF_04, sdf04Executor);
executors.put(SdfTable.SDF_13, sdf13Executor);
return executors;
}
}

View file

@ -8,6 +8,8 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.sdf.SDf01;
import ru.clearing.classes.statics.data.sdf.SDf04;
import ru.clearing.classes.statics.data.sdf.SDf13;
import ru.clearing.classes.statics.data.sdf.SDf57;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
@ -41,7 +43,7 @@ public class StatementService extends QueueConsumer implements InitializingBean
/**
* Пары sdf запросов пришедшие с модуля dbf-import
*/
private final List<Pair<StatementRequest, StatementRequest>> pairOfSdfRequest = new LinkedList<>();
private final Map<Long, Pair<StatementRequest, StatementRequest>> pairOfSdfRequest = new HashMap<>();
@Autowired
public StatementService(Consumer<String, Object> kafkaQueue,
@ -55,6 +57,8 @@ public class StatementService extends QueueConsumer implements InitializingBean
this.executorsMap = executorsMap;
this.sdfImdgs.put(SdfTable.SDF_01, imdgProvider.getImdg(IMDGDistributedNames.Map_SDf01, SDf01.class));
this.sdfImdgs.put(SdfTable.SDF_57, imdgProvider.getImdg(IMDGDistributedNames.Map_SDf57, SDf57.class));
this.sdfImdgs.put(SdfTable.SDF_04, imdgProvider.getImdg(IMDGDistributedNames.Map_SDf04, SDf04.class));
this.sdfImdgs.put(SdfTable.SDF_13, imdgProvider.getImdg(IMDGDistributedNames.Map_SDf13, SDf13.class));
}
@Override
@ -67,16 +71,15 @@ public class StatementService extends QueueConsumer implements InitializingBean
private void process(BaseRequest<StatementRequest> systemRequest) {
StatementRequest statementRequest = systemRequest.getRequestPayload();
Collection<? extends SpcexObjectBase> sdfGroup;
SdfTable table = statementRequest.getTable();
Imdg<? extends SpcexObjectBase> sdfImdg = sdfImdgs.get(table);
Optional<Pair<StatementRequest, StatementRequest>> completePairOpt = saveRequest(statementRequest);
if (completePairOpt.isPresent()) {
Optional<Long> completePairKey = saveRequest(statementRequest);
if (completePairKey.isPresent()) {
if (List.of(SdfTable.SDF_01, SdfTable.SDF_57).contains(table)) {
{
//всегда сначала обработаем sdf57
Pair<StatementRequest, StatementRequest> pair = completePairOpt.get();
Long key = completePairKey.get();
Pair<StatementRequest, StatementRequest> pair = pairOfSdfRequest.get(key);
processSdf57(pair.getSecond());
//затем sdf01
@ -84,9 +87,23 @@ public class StatementService extends QueueConsumer implements InitializingBean
//теперь можем продолжить сессию с шага 1
ContinueSessionBnRequest continueSessionBn = new ContinueSessionBnRequest();
kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN_FIRST_PART, continueSessionBn);
pairOfSdfRequest.remove(key);
}
} else if (List.of(SdfTable.SDF_08, SdfTable.SDF_13).contains(table)) {
//not implemented part; it's actually stage number 8 from any session
{
//всегда сначала обработаем sdf04
Long key = completePairKey.get();
Pair<StatementRequest, StatementRequest> pair = pairOfSdfRequest.get(key);
processSdf04(pair.getFirst());
//затем sdf01
processSdf13(pair.getSecond());
//теперь можем продолжить сессию с шага 1
ContinueSessionBnRequest continueSessionBn = new ContinueSessionBnRequest();
kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN_SECOND_PART, continueSessionBn);
pairOfSdfRequest.remove(key);
}
}
}
}
@ -99,6 +116,22 @@ public class StatementService extends QueueConsumer implements InitializingBean
service.execute(sdfGroup, statementRequest);
}
private void processSdf04(StatementRequest statementRequest) {
Imdg<SDf04> sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf04, SDf04.class);
Collection<? extends SpcexObjectBase> sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of(
"generationId", statementRequest.getGroupId()));
AbstractExecutor service = executorsMap.get(SdfTable.SDF_04);
service.execute(sdfGroup, statementRequest);
}
private void processSdf13(StatementRequest statementRequest) {
Imdg<SDf13> sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf13, SDf13.class);
Collection<? extends SpcexObjectBase> sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of(
"generationId", statementRequest.getGroupId()));
AbstractExecutor service = executorsMap.get(SdfTable.SDF_13);
service.execute(sdfGroup, statementRequest);
}
private void processSdf01(StatementRequest statementRequest) {
Imdg<SDf01> sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf01, SDf01.class);
Collection<? extends SpcexObjectBase> sdfGroup;
@ -122,25 +155,31 @@ public class StatementService extends QueueConsumer implements InitializingBean
}
}
private Optional<Pair<StatementRequest, StatementRequest>> saveRequest(StatementRequest statementRequest) {
private Optional<Long> saveRequest(StatementRequest statementRequest) {
SdfTable sdfTable = statementRequest.getTable();
if (pairOfSdfRequest.isEmpty()) {
pairOfSdfRequest.add(getPairByTableName(sdfTable, statementRequest));
Long id = imdgProvider.getImdgIdGenerator().nextId();
;
pairOfSdfRequest.put(id, getPairByTableName(sdfTable, statementRequest));
} else {
//найдем первую неполноценную пару
Optional<Pair<StatementRequest, StatementRequest>> uncompletedPairOpt = findUncompletedPairByTableName(sdfTable);
if (uncompletedPairOpt.isPresent()) {
Pair<StatementRequest, StatementRequest> uncompletedPair = uncompletedPairOpt.get();
if (uncompletedPair.getFirst() == null) {
uncompletedPair.setFirst(statementRequest);
Optional<Long> uncompletedPairId = findUncompletedPairByTableName(sdfTable);
if (uncompletedPairId.isPresent()) {
Long key = uncompletedPairId.get();
Pair<StatementRequest, StatementRequest> completed = pairOfSdfRequest.get(key);
if (completed.getFirst() == null) {
completed.setFirst(statementRequest);
} else {
uncompletedPair.setSecond(statementRequest);
completed.setSecond(statementRequest);
}
pairOfSdfRequest.put(key, completed);
//укомплектованная пара
return Optional.of(uncompletedPair);
return Optional.of(key);
} else {
//если таких нет, то просто создаем новую с одной частью
pairOfSdfRequest.add(getPairByTableName(sdfTable, statementRequest));
Long id = imdgProvider.getImdgIdGenerator().nextId();
;
pairOfSdfRequest.put(id, getPairByTableName(sdfTable, statementRequest));
}
}
return Optional.empty();
@ -151,17 +190,23 @@ public class StatementService extends QueueConsumer implements InitializingBean
switch (sdfTable) {
case SDF_01 -> pair = new Pair<>(statementRequest, null);
case SDF_57 -> pair = new Pair<>(null, statementRequest);
case SDF_04 -> pair = new Pair<>(statementRequest, null);
case SDF_13 -> pair = new Pair<>(null, statementRequest);
}
return pair;
}
private Optional<Pair<StatementRequest, StatementRequest>> findUncompletedPairByTableName(SdfTable sdfTable) {
Optional<Pair<StatementRequest, StatementRequest>> uncompletedPair = Optional.empty();
private Optional<Long> findUncompletedPairByTableName(SdfTable sdfTable) {
Optional<Long> uncompletedPair = Optional.empty();
switch (sdfTable) {
case SDF_01 ->
uncompletedPair = pairOfSdfRequest.stream().filter(pair -> pair.getFirst() != null && pair.getSecond() == null).findFirst();
uncompletedPair = pairOfSdfRequest.entrySet().stream().filter(entry -> entry.getValue().getFirst() != null && entry.getValue().getSecond() == null).findFirst().map(Map.Entry::getKey);
case SDF_57 ->
uncompletedPair = pairOfSdfRequest.stream().filter(pair -> pair.getSecond() != null && pair.getFirst() == null).findFirst();
uncompletedPair = pairOfSdfRequest.entrySet().stream().filter(entry -> entry.getValue().getFirst() == null && entry.getValue().getSecond() != null).findFirst().map(Map.Entry::getKey);
case SDF_04 ->
uncompletedPair = pairOfSdfRequest.entrySet().stream().filter(entry -> entry.getValue().getFirst() != null && entry.getValue().getSecond() == null).findFirst().map(Map.Entry::getKey);
case SDF_13 ->
uncompletedPair = pairOfSdfRequest.entrySet().stream().filter(entry -> entry.getValue().getFirst() == null && entry.getValue().getSecond() != null).findFirst().map(Map.Entry::getKey);
}
return uncompletedPair;
}

View file

@ -0,0 +1,98 @@
package ru.spcex.clearing.service.executors;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.clearing.classes.statics.data.sdf.SDf04;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.model.Result;
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.imdg.api.predicate.ImdgPredicate;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
import java.math.BigDecimal;
import java.util.Collection;
import java.util.Collections;
@Service
public class Sdf04Executor extends AbstractExecutor<SDf04> {
private final Logger log = LoggerFactory.getLogger(getClass());
private final ImdgProvider imdgProvider;
private final Imdg<Registry> registryImdg;
public Sdf04Executor(ImdgProvider imdgProvider) {
this.imdgProvider = imdgProvider;
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
}
@Override
public String exportTableName() {
return "";
}
@Override
public boolean isNeedToSendCommand() {
return false;
}
@Override
public void sendCommand(KafkaSender kafkaSender, Result result) {
}
public Result execute(Collection<SDf04> sdf, StatementRequest statementRequest) {
Result result = new Result();
Long generationIdForGroup = imdgProvider.getImdgIdGenerator().nextId();
result.setGenerationId(generationIdForGroup);
for (SDf04 sdf04 : sdf) {
//обычно мы ищем по группу sdf04, здесь как будто всегда только одна запись, todo нужно прочекать этот момент
log.debug("Process sdf04 record; sdf04.id: {}", sdf04.getId());
Collection<String> fullNames = Collections.emptyList(); //fixme что значит registry.fullName=sDf04.(sbanknam1, sbanknam2 и т.д.)
Collection<Registry> registries = selectRegistryForSDF04(sdf04.getC_acc_cred(), fullNames);
registries.forEach(registry -> unlockRegistry(registry, new BigDecimal(sdf04.getPay_val())));
}
return result;
}
protected Collection<Registry> selectRegistryForSDF04(String account, Collection<String> fullNames) {
ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
ImdgPredicate query = pb.and(
pb.and(
pb.equals("registryDesignation", RegistryDesignation.A.getKey()),
pb.equals("registryInstrumentType", RegistryInstrumentType.M.getKey()),
pb.or(pb.equals("registryUnit", RegistryUnit.F.getKey()),
pb.equals("registryUnit", RegistryUnit.B.getKey()))
),
pb.equals("account", account)
// ,
// pb.in("fullName", fullNames.toArray(new String[fullNames.size()]))
);
Collection<Registry> result = registryImdg.getCollectionObjectsByPredicate(query);
log.trace("Selected {} registry's by sql: {}", result.size(), query);
return result;
}
boolean unlockRegistry(Registry registry, BigDecimal value) {
if (RegistryUnit.F.equalsByKey(registry.getRegistryUnit())) {
if (registry.getBalance() == null) registry.setBalance(BigDecimal.ZERO);
registry.setBalance(registry.getBalance().add(value));
return true;
} else if (RegistryUnit.B.equalsByKey(registry.getRegistryUnit())) {
if (registry.getBalance() == null) registry.setBalance(BigDecimal.ZERO);
registry.setBalance(registry.getBalance().subtract(value));
return true;
} else {
log.warn("For registry {} registryUnit={} unlock operation not implemented.",
registry.getId(), registry.getRegistryUnit());
return false;
}
}
}

View file

@ -0,0 +1,105 @@
package ru.spcex.clearing.service.executors;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.clearing.classes.statics.data.sdf.SDf12;
import ru.clearing.classes.statics.data.sdf.SDf13;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.model.Result;
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.imdg.api.predicate.ImdgPredicate;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
import java.math.BigDecimal;
import java.util.Collection;
@Service
public class Sdf13Executor extends AbstractExecutor<SDf13> {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<SDf12> sdf12Imdg;
private final ImdgProvider imdgProvider;
private final Imdg<Registry> registryImdg;
public Sdf13Executor(ImdgProvider imdgProvider) {
this.imdgProvider = imdgProvider;
this.sdf12Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf12, SDf12.class);
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
}
@Override
public String exportTableName() {
return "";
}
@Override
public boolean isNeedToSendCommand() {
return false;
}
@Override
public void sendCommand(KafkaSender kafkaSender, Result result) {
}
public Result execute(Collection<SDf13> sdf, StatementRequest statementRequest) {
Result result = new Result();
Long generationIdForGroup = imdgProvider.getImdgIdGenerator().nextId();
result.setGenerationId(generationIdForGroup);
for (SDf13 sdf13 : sdf) {
//обычно мы ищем по группу sdf04, здесь как будто всегда только одна запись, todo нужно прочекать этот момент
SDf12 sDf12 = selectSdf12bySdf13(sdf13);
log.debug("Process sdf13 record; sdf13.id: {}", sdf13.getId());
Collection<Registry> registries = selectRegistryForSDF12(sDf12.getDepoCodeSender(), sDf12.getSecurityCode());
registries.forEach(registry -> unlockRegistry(registry, new BigDecimal(sDf12.getQuantity())));
}
return result;
}
protected SDf12 selectSdf12bySdf13(SDf13 sDf13) {
SDf12 sdf12 = sdf12Imdg.getSingleObjectByID(sDf13.getId());
log.trace("Selected sdf12 by sdf13");
return sdf12;
}
protected Collection<Registry> selectRegistryForSDF12(String account, String securityCode) {
ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
ImdgPredicate query = pb.and(
pb.and(
pb.equals("registryDesignation", RegistryDesignation.A.getKey()),
pb.equals("registryInstrumentType", RegistryInstrumentType.M.getKey()),
pb.or(pb.equals("registryUnit", RegistryUnit.F.getKey()),
pb.equals("registryUnit", RegistryUnit.B.getKey()))
),
pb.equals("account", account),
pb.equals("securityCode", securityCode)
);
Collection<Registry> result = registryImdg.getCollectionObjectsByPredicate(query);
log.trace("Selected {} registry's by sql: {}", result.size(), query);
return result;
}
boolean unlockRegistry(Registry registry, BigDecimal value) {
if (RegistryUnit.F.equalsByKey(registry.getRegistryUnit())) {
if (registry.getBalance() == null) registry.setBalance(BigDecimal.ZERO);
registry.setBalance(registry.getBalance().add(value));
return true;
} else if (RegistryUnit.B.equalsByKey(registry.getRegistryUnit())) {
if (registry.getBalance() == null) registry.setBalance(BigDecimal.ZERO);
registry.setBalance(registry.getBalance().subtract(value));
return true;
} else {
log.warn("For registry {} registryUnit={} unlock operation not implemented.",
registry.getId(), registry.getRegistryUnit());
return false;
}
}
}