fix statementService
This commit is contained in:
parent
cec531d04a
commit
389eb7d9fd
2 changed files with 91 additions and 8 deletions
|
|
@ -15,6 +15,7 @@ import ru.spcex.clearing.platform.messaging.domain.Consts;
|
|||
import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdf01Request;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdfRequestPart;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.ContinueSessionBnRequest;
|
||||
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||
import ru.spcex.clearing.service.executors.AbstractExecutor;
|
||||
|
|
@ -23,11 +24,9 @@ import ru.spcex.platform.classes.base.SpcexObjectBase;
|
|||
import ru.spcex.platform.enumeration.SdfTable;
|
||||
import ru.spcex.platform.imdg.api.Imdg;
|
||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||
import ru.spcex.platform.utils.collection.Pair;
|
||||
|
||||
import java.util.Collection;
|
||||
import java.util.EnumMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.*;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
@Service
|
||||
|
|
@ -39,6 +38,11 @@ public class StatementService extends QueueConsumer implements InitializingBean
|
|||
private final Map<SdfTable, Imdg<? extends SpcexObjectBase>> sdfImdgs;
|
||||
private final Map<SdfTable, AbstractExecutor<?>> executorsMap;
|
||||
|
||||
/**
|
||||
* Пары sdf запросов пришедшие с модуля dbf-import
|
||||
*/
|
||||
private final List<Pair<StatementRequest, StatementRequest>> pairOfSdfRequest = new LinkedList<>();
|
||||
|
||||
@Autowired
|
||||
public StatementService(Consumer<String, Object> kafkaQueue,
|
||||
ImdgProvider imdgProvider,
|
||||
|
|
@ -66,17 +70,50 @@ public class StatementService extends QueueConsumer implements InitializingBean
|
|||
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()) {
|
||||
if (List.of(SdfTable.SDF_01, SdfTable.SDF_57).contains(table)) {
|
||||
{
|
||||
//всегда сначала обработаем sdf57
|
||||
Pair<StatementRequest, StatementRequest> pair = completePairOpt.get();
|
||||
processSdf57(pair.getSecond());
|
||||
|
||||
//затем sdf01
|
||||
processSdf01(pair.getFirst());
|
||||
//теперь можем продолжить сессию с шага 1
|
||||
ContinueSessionBnRequest continueSessionBn = new ContinueSessionBnRequest();
|
||||
kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN_FIRST_PART, continueSessionBn);
|
||||
}
|
||||
} else if (List.of(SdfTable.SDF_08, SdfTable.SDF_13).contains(table)) {
|
||||
//not implemented part; it's actually stage number 8 from any session
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void processSdf57(StatementRequest statementRequest) {
|
||||
Imdg<SDf57> sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf57, SDf57.class);
|
||||
Collection<? extends SpcexObjectBase> sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of(
|
||||
"generationId", statementRequest.getGroupId()));
|
||||
AbstractExecutor service = executorsMap.get(SdfTable.SDF_57);
|
||||
service.execute(sdfGroup, statementRequest);
|
||||
}
|
||||
|
||||
private void processSdf01(StatementRequest statementRequest) {
|
||||
Imdg<SDf01> sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf01, SDf01.class);
|
||||
Collection<? extends SpcexObjectBase> sdfGroup;
|
||||
if (statementRequest.getAccountCreationResults().size() == 0) {
|
||||
sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of("generationId", statementRequest.getGroupId()));
|
||||
sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of(
|
||||
"generationId", statementRequest.getGroupId()));
|
||||
} else {
|
||||
//todo могут ли тут реквесты бегать по кругу, если да, то у нас проблемы
|
||||
sdfGroup = statementRequest.getAccountCreationResults()
|
||||
.stream()
|
||||
.filter(part -> part.getErrorCode() == null) //fixme эти случае должны попадать в ошибочный sdf02
|
||||
.map(part -> sdfImdg.getSingleObjectByID(part.getSdfId()))
|
||||
.collect(Collectors.toList());
|
||||
}
|
||||
|
||||
AbstractExecutor service = executorsMap.get(table);
|
||||
AbstractExecutor service = executorsMap.get(SdfTable.SDF_01);
|
||||
Result res = service.execute(sdfGroup, statementRequest);
|
||||
if (res.getAccountRequests().size() != 0) {
|
||||
kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF01, createAccountsRequest(statementRequest.getGroupId(), res.getAccountRequests()));
|
||||
|
|
@ -85,6 +122,50 @@ public class StatementService extends QueueConsumer implements InitializingBean
|
|||
}
|
||||
}
|
||||
|
||||
private Optional<Pair<StatementRequest, StatementRequest>> saveRequest(StatementRequest statementRequest) {
|
||||
SdfTable sdfTable = statementRequest.getTable();
|
||||
if (pairOfSdfRequest.isEmpty()) {
|
||||
pairOfSdfRequest.add(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);
|
||||
} else {
|
||||
uncompletedPair.setSecond(statementRequest);
|
||||
}
|
||||
//укомплектованная пара
|
||||
return Optional.of(uncompletedPair);
|
||||
} else {
|
||||
//если таких нет, то просто создаем новую с одной частью
|
||||
pairOfSdfRequest.add(getPairByTableName(sdfTable, statementRequest));
|
||||
}
|
||||
}
|
||||
return Optional.empty();
|
||||
}
|
||||
|
||||
private Pair<StatementRequest, StatementRequest> getPairByTableName(SdfTable sdfTable, StatementRequest statementRequest) {
|
||||
Pair<StatementRequest, StatementRequest> pair = null;
|
||||
switch (sdfTable) {
|
||||
case SDF_01 -> pair = new Pair<>(statementRequest, null);
|
||||
case SDF_57 -> pair = new Pair<>(null, statementRequest);
|
||||
}
|
||||
return pair;
|
||||
}
|
||||
|
||||
private Optional<Pair<StatementRequest, StatementRequest>> findUncompletedPairByTableName(SdfTable sdfTable) {
|
||||
Optional<Pair<StatementRequest, StatementRequest>> uncompletedPair = Optional.empty();
|
||||
switch (sdfTable) {
|
||||
case SDF_01 ->
|
||||
uncompletedPair = pairOfSdfRequest.stream().filter(pair -> pair.getFirst() != null && pair.getSecond() == null).findFirst();
|
||||
case SDF_57 ->
|
||||
uncompletedPair = pairOfSdfRequest.stream().filter(pair -> pair.getSecond() != null && pair.getFirst() == null).findFirst();
|
||||
}
|
||||
return uncompletedPair;
|
||||
}
|
||||
|
||||
private AccountSdf01Request createAccountsRequest(Long sdf01GroupingId, List<AccountSdfRequestPart> accountRequests) {
|
||||
AccountSdf01Request r = new AccountSdf01Request();
|
||||
r.setGroupingSdf01Id(sdf01GroupingId);
|
||||
|
|
|
|||
|
|
@ -3,7 +3,9 @@ package ru.spcex.platform.enumeration;
|
|||
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
||||
|
||||
public enum SdfTable implements IEnumKey {
|
||||
SDF_01("SDF_01"), SDF_57("SDF_57"), SDF_09("SDF_09"), SDF_16("SDF_16");
|
||||
SDF_01("SDF_01"), SDF_04("SDF_04"), SDF_08("SDF_08"),
|
||||
SDF_13("SDF_13"), SDF_57("SDF_57"), SDF_09("SDF_09"),
|
||||
SDF_16("SDF_16");
|
||||
|
||||
SdfTable(String key) {
|
||||
this.key = key;
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue