sdf08 sdf21

This commit is contained in:
ialbert 2023-08-08 15:45:01 +03:00
parent fb0b57bffc
commit 0383068815
2 changed files with 176 additions and 3 deletions

View file

@ -38,6 +38,7 @@ public class StatementService extends QueueConsumer implements InitializingBean
private final Map<SdfTable, Imdg<? extends SpcexObjectBase>> sdfImdgs;
private final Map<SdfTable, AbstractExecutor<?>> executorsMap;
private final Reviser reviser;
private final StatementServiceV2 stmtSrvV2;
/**
* Пары sdf запросов пришедшие с модуля dbf-import
@ -48,10 +49,11 @@ public class StatementService extends QueueConsumer implements InitializingBean
public StatementService(Consumer<String, Object> kafkaQueue,
ImdgProvider imdgProvider,
KafkaSender kafkaSender,
@Qualifier("sdfExecutors") Map<SdfTable, AbstractExecutor<?>> executorsMap, Reviser reviser) {
@Qualifier("sdfExecutors") Map<SdfTable, AbstractExecutor<?>> executorsMap, Reviser reviser, StatementServiceV2 stmtSrvV2) {
super(kafkaQueue);
this.imdgProvider = imdgProvider;
this.reviser = reviser;
this.stmtSrvV2 = stmtSrvV2;
this.sdfImdgs = new EnumMap<>(SdfTable.class);
this.kafkaSender = kafkaSender;
this.executorsMap = executorsMap;
@ -68,7 +70,9 @@ public class StatementService extends QueueConsumer implements InitializingBean
callback(StatementRequest.class)
.setConsumer(systemRequest -> {
StatementRequest payload = systemRequest.getRequestPayload();
if (payload.isContinueSdf()) {
if (Arrays.asList(SdfTable.SDF_08, SdfTable.SDF_21).contains(payload.getTable())) {
stmtSrvV2.processReq(systemRequest);
} else if (payload.isContinueSdf()) {
processAccountAnswer(systemRequest);
} else {
processPaired(systemRequest);
@ -382,7 +386,7 @@ public class StatementService extends QueueConsumer implements InitializingBean
return r;
}
private AccountSdf01Request createAccountsRequest(Long sdf01GroupingId, List<AccountSdfRequestPart> accountRequests, Long childGenerationId) {
public static AccountSdf01Request createAccountsRequest(Long sdf01GroupingId, List<AccountSdfRequestPart> accountRequests, Long childGenerationId) {
AccountSdf01Request r = new AccountSdf01Request();
r.setGroupingSdf01Id(sdf01GroupingId);
r.setGroupingSdf02Id(childGenerationId);

View file

@ -0,0 +1,169 @@
package ru.spcex.clearing.service;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;
import ru.clearing.classes.statics.data.sdf.SDf08;
import ru.clearing.classes.statics.data.sdf.SDf21;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
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.balance.StatementRequest;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.executors.Sdf08Executor;
import ru.spcex.clearing.service.executors.Sdf21Executor;
import ru.spcex.clearing.service.model.Result;
import ru.spcex.platform.enumeration.SdfTable;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import java.util.*;
import java.util.stream.Collectors;
@Component
public class StatementServiceV2 {
private Logger log = LoggerFactory.getLogger(getClass());
private final ImdgProvider imdgProvider;
private final KafkaSender kafkaSender;
private final List<StatementRequest> statementRequests = new LinkedList<>();
private final Sdf08Executor sdf08Executor;
private final Sdf21Executor sdf21Executor;
public StatementServiceV2(ImdgProvider imdgProvider, KafkaSender kafkaSender, Sdf08Executor sdf08Executor, Sdf21Executor sdf21Executor) {
this.imdgProvider = imdgProvider;
this.kafkaSender = kafkaSender;
this.sdf08Executor = sdf08Executor;
this.sdf21Executor = sdf21Executor;
}
public void processReq(BaseRequest<StatementRequest> systemRequest) {
SdfTable table = systemRequest.getRequestPayload().getTable();
Long groupId = systemRequest.getRequestPayload().getGroupId();
log.debug("received stmtReq.SdfTable={}, groupId={}", table, groupId);
Optional<SdfGroup> sdfGroup = SdfGroup.groupByTable(table);
if (sdfGroup.isEmpty()) {
log.debug("table {} is not paired with any other SDF", table);
//здесь switch case одиночные методы
return;
}
Collection<StatementRequest> fullGroup = getFullGroup(systemRequest.getRequestPayload(), sdfGroup.get());
if (fullGroup.isEmpty()) {
log.debug("didn't find full sdf set for stmtReq.SdfTable={}, groupId={}. Adding to cache", table, groupId);
statementRequests.add(systemRequest.getRequestPayload());
return;
}
log.debug("find full set of SDF requests: {}", fullGroup.stream()
.map(stReq -> stReq.getTable().getKey() + " generationId=" + stReq.getGroupId())
.collect(Collectors.joining(",", "[", "]")));
if (sdfGroup.get() == SdfGroup.Sdf08And21) {
processSdf08And21(find(SdfTable.SDF_08, fullGroup), find(SdfTable.SDF_21, fullGroup));
} else {
throw new IllegalStateException("not implemented");
}
}
private static StatementRequest find(SdfTable table, Collection<StatementRequest> reqs) {
Optional<StatementRequest> first = reqs.stream().filter(r -> r.getTable().equals(table)).findFirst();
if (first.isEmpty()) {
throw new IllegalStateException("couldn't find + " + table + " - this should have happened ");
}
return first.get();
}
/**
* fromAccService передается когда пришел ответ от account-service
* в этом случае: по key находим пару в которой сохранен sdf57 запрос и частично выполненный sdf01
* вместо старого sdf01 запроса выполняем новый пришедший от account-service
*/
private void processSdf08And21(StatementRequest sdf08, StatementRequest sdf21) {
Result sdf08Res = processSdf08(sdf08);
if (sdf08Res.getAccountRequests().size() > 0) {
statementRequests.remove(sdf08);
log.info("sdf08 execution wasn't complete, waiting for an answer from account-service");
return;
}
//затем sdf21
processSdf21(sdf21);
//fixme ревизия для бумаг reviser.doRevise(pair.getFirst().getGroupId());
log.info("pair sdf08/sdf21 processed successfully");
statementRequests.removeAll(Arrays.asList(sdf08, sdf21));
}
private Result processSdf08(StatementRequest statementRequest) {
Imdg<SDf08> sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf08, SDf08.class);
Collection<SDf08> sdfGroup;
if (statementRequest.getAccountCreationResults().size() == 0) {
sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of(
"generationId", statementRequest.getGroupId()));
} else {
sdfGroup = statementRequest.getAccountCreationResults()
.stream()
.filter(part -> part.getErrorCode() == null)
.map(part -> sdfImdg.getSingleObjectByID(part.getSdfId()))
.collect(Collectors.toList());
}
Result res = sdf08Executor.execute(sdfGroup, statementRequest);
if (res.getAccountRequests().size() != 0) {
AccountSdf01Request createAccsReq = StatementService.createAccountsRequest(statementRequest.getGroupId(),
res.getAccountRequests(),
res.getGenerationId());
kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF08, createAccsReq);
} else if (sdf08Executor.isNeedToSendCommand()) {
sdf08Executor.sendCommand(kafkaSender, res);
}
return res;
}
private void processSdf21(StatementRequest statementRequest) {
Imdg<SDf21> sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf21, SDf21.class);
Collection<SDf21> sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of(
"generationId", statementRequest.getGroupId()));
Result res = sdf21Executor.execute(sdfGroup, statementRequest);
}
private Collection<StatementRequest> getFullGroup(StatementRequest stmtReq, SdfGroup group) {
List<StatementRequest> reqs = new ArrayList<>();
log.debug("trying to find request pair for {}. cashed reqs {}", stmtReq.getTable(),
statementRequests.stream()
.map(StatementRequest::getTable)
.map(Object::toString)
.collect(Collectors.joining(",", "[", "]")));
for (StatementRequest stmtReqOld : statementRequests) { //проходим с головы (с самых старых)
boolean matchGroup = group.getGroup().contains(stmtReqOld.getTable());
boolean notIncludedThisTableAlready = reqs.stream().map(StatementRequest::getTable).noneMatch(t -> t.equals(stmtReqOld.getTable()));
boolean isNotTheSameAsNewReq = Objects.equals(stmtReqOld.getTable(), stmtReq.getTable());
if (matchGroup && notIncludedThisTableAlready && isNotTheSameAsNewReq) {
reqs.add(stmtReqOld);
}
}
reqs.add(stmtReq);
if (reqs.size() == group.getGroup().size()) {
return reqs;
} else {
return Collections.emptyList();
}
}
private enum SdfGroup {
Sdf01And57(SdfTable.SDF_01, SdfTable.SDF_57),
Sdf08And21(SdfTable.SDF_08, SdfTable.SDF_21);
private final Collection<SdfTable> group;
SdfGroup(SdfTable... group) {
this.group = List.of(group);
}
public Collection<SdfTable> getGroup() {
return group;
}
public static Optional<SdfGroup> groupByTable(SdfTable t) {
return Arrays.stream(SdfGroup.values())
.filter(g -> g.getGroup().contains(t))
.findFirst();
}
}
}