SDF группы для сессий

избавился от StatementService.java
This commit is contained in:
ialbert 2023-09-05 19:14:17 +03:00
parent 0eadc2f7e3
commit 21df705094
6 changed files with 266 additions and 16 deletions

View file

@ -29,6 +29,7 @@ import ru.spcex.clearing.service.executors.Sdf10Executor;
import ru.spcex.clearing.service.payment.PaymentInstructionOutboundService; import ru.spcex.clearing.service.payment.PaymentInstructionOutboundService;
import ru.spcex.clearing.session.stage.*; import ru.spcex.clearing.session.stage.*;
import ru.spcex.clearing.session.stage.impl.BalanceRevise; import ru.spcex.clearing.session.stage.impl.BalanceRevise;
import ru.spcex.clearing.statement.StatementServiceV2;
import ru.spcex.platform.enumeration.Task; import ru.spcex.platform.enumeration.Task;
import ru.spcex.platform.utils.enumeration.IMessageResolver; import ru.spcex.platform.utils.enumeration.IMessageResolver;
import ru.spcex.platform.utils.error.ValidationException; import ru.spcex.platform.utils.error.ValidationException;
@ -53,6 +54,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
private final Sdf10Executor sdf10Executor; private final Sdf10Executor sdf10Executor;
private final BalanceRevise balanceRevise; private final BalanceRevise balanceRevise;
private final Sdf05Sender sdf05Sender; private final Sdf05Sender sdf05Sender;
private final StatementServiceV2 statementService;
private final PaymentInstructionOutboundService pmtOutboundService; private final PaymentInstructionOutboundService pmtOutboundService;
@Autowired @Autowired
@ -64,7 +66,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
SecondaryAuctionT0Session secondaryAuctionT0Session, SecondaryAuctionT0Session secondaryAuctionT0Session,
PrimaryAuctionB0Session primaryAuctionB0Session, PrimaryAuctionT0Session primaryAuctionT0Session, IntermediateMkrSession intermediateMkrSession, FinalMkrSession finalMkrSession, ReturnDepositSession returnDepositSession, SessionManager sessionManager, PrimaryAuctionB0Session primaryAuctionB0Session, PrimaryAuctionT0Session primaryAuctionT0Session, IntermediateMkrSession intermediateMkrSession, FinalMkrSession finalMkrSession, ReturnDepositSession returnDepositSession, SessionManager sessionManager,
Sdf06Executor sdf06Executor, Sdf06Executor sdf06Executor,
Sdf10Executor sdf10Executor, BalanceRevise balanceRevise, Sdf05Sender sdf05Sender, PaymentInstructionOutboundService pmtOutboundService) { Sdf10Executor sdf10Executor, BalanceRevise balanceRevise, Sdf05Sender sdf05Sender, StatementServiceV2 statementService, PaymentInstructionOutboundService pmtOutboundService) {
super(kafkaQueue, kafkaResponseQueue); super(kafkaQueue, kafkaResponseQueue);
this.errorResolver = errorResolver; this.errorResolver = errorResolver;
this.clearingService = clearingService; this.clearingService = clearingService;
@ -81,6 +83,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
this.sdf10Executor = sdf10Executor; this.sdf10Executor = sdf10Executor;
this.balanceRevise = balanceRevise; this.balanceRevise = balanceRevise;
this.sdf05Sender = sdf05Sender; this.sdf05Sender = sdf05Sender;
this.statementService = statementService;
this.pmtOutboundService = pmtOutboundService; this.pmtOutboundService = pmtOutboundService;
} }
@ -116,6 +119,9 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
callback(PIClearingOutbondActionNewRequest.class) callback(PIClearingOutbondActionNewRequest.class)
.setFunction(pmtOutboundService::sendOutBoundPayment) .setFunction(pmtOutboundService::sendOutBoundPayment)
.forDestination(Consts.PAYMENT_INSTRUCTION_CLEARING_OUTBOUND_ACTION, callbacks::put); .forDestination(Consts.PAYMENT_INSTRUCTION_CLEARING_OUTBOUND_ACTION, callbacks::put);
callback(StatementRequest.class)
.setConsumer(statementService::processReq)
.forDestination(Consts.STATEMENT_PROCESS, callbacks::put);
callback(SessionContinueEvent.class) callback(SessionContinueEvent.class)
.setConsumer(req -> { .setConsumer(req -> {
primaryAuctionBnSession.continueSession(req); primaryAuctionBnSession.continueSession(req);

View file

@ -4,9 +4,7 @@ import org.apache.kafka.clients.consumer.Consumer;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.sdf.*; import ru.clearing.classes.statics.data.sdf.*;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
@ -33,7 +31,8 @@ import ru.spcex.platform.utils.collection.Pair;
import java.util.*; import java.util.*;
import java.util.stream.Collectors; import java.util.stream.Collectors;
@Service //@Service
@Deprecated
public class StatementService extends QueueConsumer implements InitializingBean { public class StatementService extends QueueConsumer implements InitializingBean {
private final Logger log = LoggerFactory.getLogger(getClass()); private final Logger log = LoggerFactory.getLogger(getClass());
@ -49,7 +48,7 @@ public class StatementService extends QueueConsumer implements InitializingBean
*/ */
private final Map<Long, Pair<StatementRequest, StatementRequest>> pairOfSdfRequest = new HashMap<>(); private final Map<Long, Pair<StatementRequest, StatementRequest>> pairOfSdfRequest = new HashMap<>();
@Autowired // @Autowired
public StatementService(Consumer<String, Object> kafkaQueue, public StatementService(Consumer<String, Object> kafkaQueue,
ImdgProvider imdgProvider, ImdgProvider imdgProvider,
KafkaSender kafkaSender, KafkaSender kafkaSender,
@ -75,11 +74,28 @@ public class StatementService extends QueueConsumer implements InitializingBean
.setConsumer(systemRequest -> { .setConsumer(systemRequest -> {
StatementRequest payload = systemRequest.getRequestPayload(); StatementRequest payload = systemRequest.getRequestPayload();
if (payload.getTable() == null if (payload.getTable() == null
|| !Arrays.asList(SdfTable.SDF_01, SdfTable.SDF_04, SdfTable.SDF_57, SdfTable.SDF_08, SdfTable.SDF_13, SdfTable.SDF_21, SdfTable.SDF_10, SdfTable.SDF_55, SdfTable.SDF_20).contains(payload.getTable())) { || !Arrays.asList(SdfTable.SDF_01,
SdfTable.SDF_04,
SdfTable.SDF_57,
SdfTable.SDF_08,
SdfTable.SDF_13,
SdfTable.SDF_21,
SdfTable.SDF_10,
SdfTable.SDF_55,
SdfTable.SDF_20).contains(payload.getTable())) {
log.debug("StatementService: skipping table {}", payload.getTable()); log.debug("StatementService: skipping table {}", payload.getTable());
return; return;
} }
if (Arrays.asList(SdfTable.SDF_08, SdfTable.SDF_21, SdfTable.SDF_10, SdfTable.SDF_55, SdfTable.SDF_20).contains(payload.getTable())) { if (Arrays.asList(
SdfTable.SDF_01,
SdfTable.SDF_04,
SdfTable.SDF_08,
SdfTable.SDF_10,
SdfTable.SDF_13,
SdfTable.SDF_20,
SdfTable.SDF_21,
SdfTable.SDF_55,
SdfTable.SDF_57).contains(payload.getTable())) {
stmtSrvV2.processReq(systemRequest); stmtSrvV2.processReq(systemRequest);
} else if (payload.isContinueSdf()) { } else if (payload.isContinueSdf()) {
processAccountAnswer(systemRequest); processAccountAnswer(systemRequest);

View file

@ -13,7 +13,12 @@ public enum SdfGroup {
//------ session groups ------ //------ session groups ------
session_Triple(SdfTable.SDF_04, SdfTable.SDF_01, SdfTable.SDF_57), session_Triple(SdfTable.SDF_04, SdfTable.SDF_01, SdfTable.SDF_57),
session_Quad(SdfTable.SDF_04, SdfTable.SDF_13, SdfTable.SDF_01, SdfTable.SDF_57), session_Six(SdfTable.SDF_04,
SdfTable.SDF_13,
SdfTable.SDF_08,
SdfTable.SDF_21,
SdfTable.SDF_01,
SdfTable.SDF_57),
; ;
private final Collection<SdfTable> group; private final Collection<SdfTable> group;

View file

@ -0,0 +1,61 @@
package ru.spcex.clearing.statement;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import ru.clearing.classes.statics.data.misc.Session;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.session.stage.TaskType;
import ru.spcex.platform.classes.base.SpcexObjectBase;
import ru.spcex.platform.enumeration.SdfTable;
import ru.spcex.platform.enumeration.Section;
import ru.spcex.platform.enumeration.SessionStatus;
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 ru.spcex.platform.utils.enumeration.IEnumKey;
import ru.spcex.platform.utils.text.TextUtil;
import java.util.Collection;
import java.util.Optional;
@Component
public class SdfGroupManager {
private final Imdg<Session> sessImdg;
@Autowired
public SdfGroupManager(ImdgProvider imdgProvider) {
sessImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class);
}
public Optional<SdfGroup> getGroup(SdfTable table) {
ImdgPredicateBuilder pb = sessImdg.predicateBuilder();
ImdgPredicate prdct = pb.and(
pb.equals("workflowStatus", SessionStatus.ACTV.getKey()),
pb.equals("sessionStatus", TaskType.FormingPaymentInstruction.getKey())
);
Collection<Session> sessions = sessImdg.getCollectionObjectsByPredicate(prdct);
if (sessions.size() > 1) {
throw new IllegalStateException("more than one active session found by predicate "
+ prdct
+ " ids " + sessions.stream()
.map(SpcexObjectBase::getId)
.map(Object::toString)
.collect(TextUtil.join));
}
if (sessions.size() == 0) {
return SdfGroup.groupByTable(table);
} else {
Session session = sessions.iterator().next();
Section section = IEnumKey.getEnumByKey(Section.class, session.getSection());
if (section == null) {
throw new IllegalStateException("couldn't determine section of session.id=" + session.getId());
}
if (Section.MKR.equals(section)) {
return Optional.of(SdfGroup.session_Triple);
} else if (Section.FOND.equals(section)) {
return Optional.of(SdfGroup.session_Six);
} else throw new IllegalStateException("unknown section '" + section + "' of session.id=" + session.getId());
}
}
}

View file

@ -3,20 +3,25 @@ package ru.spcex.clearing.statement;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import ru.clearing.classes.statics.data.sdf.SDf08; import ru.clearing.classes.statics.data.sdf.*;
import ru.clearing.classes.statics.data.sdf.SDf21;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.Consts; 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.AccountSdf01Request;
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest; import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SessionContinueEvent;
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationNewRequest;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.StatementService; import ru.spcex.clearing.service.StatementService;
import ru.spcex.clearing.service.executors.*; import ru.spcex.clearing.service.executors.*;
import ru.spcex.clearing.service.model.Result; import ru.spcex.clearing.service.model.Result;
import ru.spcex.platform.classes.base.SpcexObjectBase;
import ru.spcex.platform.enumeration.ObjectType;
import ru.spcex.platform.enumeration.Priority;
import ru.spcex.platform.enumeration.SdfTable; import ru.spcex.platform.enumeration.SdfTable;
import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider; import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.utils.text.TextUtil;
import java.util.*; import java.util.*;
import java.util.function.Predicate; import java.util.function.Predicate;
@ -33,6 +38,12 @@ public class StatementServiceV2 {
private final Sdf10Executor sdf10Executor; private final Sdf10Executor sdf10Executor;
private final Sdf55Executor sdf55Executor; private final Sdf55Executor sdf55Executor;
private final Sdf20Executor sdf20Executor; private final Sdf20Executor sdf20Executor;
private final Sdf01Executor sdf01Executor;
private final Sdf57Executor sdf57Executor;
private final Sdf04Executor sdf04Executor;
private final Sdf13Executor sdf13Executor;
private final Reviser reviser;
private final SdfGroupManager grpMng;
public StatementServiceV2(ImdgProvider imdgProvider, public StatementServiceV2(ImdgProvider imdgProvider,
KafkaSender kafkaSender, KafkaSender kafkaSender,
@ -40,7 +51,7 @@ public class StatementServiceV2 {
Sdf21Executor sdf21Executor, Sdf21Executor sdf21Executor,
Sdf10Executor sdf10Executor, Sdf10Executor sdf10Executor,
Sdf55Executor sdf55Executor, Sdf55Executor sdf55Executor,
Sdf20Executor sdf20Executor) { Sdf20Executor sdf20Executor, Sdf01Executor sdf01Executor, Sdf57Executor sdf57Executor, Sdf04Executor sdf04Executor, Sdf13Executor sdf13Executor, Reviser reviser, SdfGroupManager grpMng) {
this.imdgProvider = imdgProvider; this.imdgProvider = imdgProvider;
this.kafkaSender = kafkaSender; this.kafkaSender = kafkaSender;
this.sdf08Executor = sdf08Executor; this.sdf08Executor = sdf08Executor;
@ -48,35 +59,78 @@ public class StatementServiceV2 {
this.sdf10Executor = sdf10Executor; this.sdf10Executor = sdf10Executor;
this.sdf55Executor = sdf55Executor; this.sdf55Executor = sdf55Executor;
this.sdf20Executor = sdf20Executor; this.sdf20Executor = sdf20Executor;
this.sdf01Executor = sdf01Executor;
this.sdf57Executor = sdf57Executor;
this.sdf04Executor = sdf04Executor;
this.sdf13Executor = sdf13Executor;
this.reviser = reviser;
this.grpMng = grpMng;
} }
public void processReq(BaseRequest<StatementRequest> systemRequest) { public void processReq(BaseRequest<StatementRequest> systemRequest) {
SdfTable table = systemRequest.getRequestPayload().getTable(); SdfTable table = systemRequest.getRequestPayload().getTable();
Long groupId = systemRequest.getRequestPayload().getGroupId(); Long groupId = systemRequest.getRequestPayload().getGroupId();
if (groupId == null || table == null || !Arrays.asList(
SdfTable.SDF_01,
SdfTable.SDF_04,
SdfTable.SDF_08,
SdfTable.SDF_10,
SdfTable.SDF_13,
SdfTable.SDF_20,
SdfTable.SDF_21,
SdfTable.SDF_55,
SdfTable.SDF_57).contains(table)) {
log.info("skipping table {} groupId {}", table, groupId);
return;
}
log.debug("received stmtReq.SdfTable={}, groupId={}", table, groupId); log.debug("received stmtReq.SdfTable={}, groupId={}", table, groupId);
Optional<SdfGroup> sdfGroup = SdfGroup.groupByTable(table); Optional<SdfGroup> sdfGroup = grpMng.getGroup(table);
if (sdfGroup.isEmpty()) { if (sdfGroup.isEmpty()) {
log.debug("table {} is not paired with any other SDF", table); log.debug("table {} is not paired with any other SDF", table);
//здесь switch case одиночные методы //здесь switch case одиночные методы
switch (table) { switch (table) {
case SDF_04 -> processSdf04(systemRequest.getRequestPayload());
case SDF_10 -> sdf10Executor.execute(systemRequest); case SDF_10 -> sdf10Executor.execute(systemRequest);
case SDF_55 -> sdf55Executor.execute(systemRequest); case SDF_13 -> processSdf13(systemRequest.getRequestPayload());
case SDF_20 -> sdf20Executor.execute(systemRequest); case SDF_20 -> sdf20Executor.execute(systemRequest);
case SDF_55 -> sdf55Executor.execute(systemRequest);
default -> log.error("unknown table {}", table); default -> log.error("unknown table {}", table);
} }
return; return;
} }
Collection<StatementRequest> fullGroup = getFullGroup(systemRequest.getRequestPayload(), sdfGroup.get()); log.debug("matched table {} with group {} {}", table, sdfGroup.get(),
sdfGroup.get()
.getGroup()
.stream()
.map(SdfTable::getKey)
.collect(TextUtil.join)
);
//здесь мэтчинг запросов которые пришли раньше и текущего в одну группу
//если группа нашлась полная - значит обрабатываем
List<StatementRequest> fullGroup = getFullGroup(systemRequest.getRequestPayload(), sdfGroup.get());
statementRequests.add(systemRequest.getRequestPayload()); statementRequests.add(systemRequest.getRequestPayload());
if (fullGroup.isEmpty()) { if (fullGroup.isEmpty()) {
log.debug("didn't find full sdf set for stmtReq.SdfTable={}, groupId={}. Adding to cache", table, groupId); log.debug("didn't find full sdf set for stmtReq.SdfTable={}, groupId={}. Adding to cache", table, groupId);
log.debug("current SDF cache: {}", statementRequests.stream()
.map(sr -> sr.getTable() + " groupId " + sr.getGroupId())
.collect(TextUtil.join));
return; return;
} }
log.debug("find full set of SDF requests: {}", fullGroup.stream() log.debug("find full set of SDF requests: {}", fullGroup.stream()
.map(stReq -> stReq.getTable().getKey() + " generationId=" + stReq.getGroupId()) .map(stReq -> stReq.getTable().getKey() + " generationId=" + stReq.getGroupId())
.collect(Collectors.joining(",", "[", "]"))); .collect(TextUtil.join));
if (sdfGroup.get() == SdfGroup.Sdf08And21) { if (sdfGroup.get() == SdfGroup.Sdf08And21) {
processSdf08And21(find(SdfTable.SDF_08, fullGroup), find(SdfTable.SDF_21, fullGroup)); processSdf08And21(find(SdfTable.SDF_08, fullGroup), find(SdfTable.SDF_21, fullGroup));
} else if (sdfGroup.get() == SdfGroup.Sdf01And57) {
processSdf01And57(find(SdfTable.SDF_01, fullGroup), find(SdfTable.SDF_57, fullGroup));
} else if (sdfGroup.get() == SdfGroup.session_Triple) {
processSdf04(find(SdfTable.SDF_04, fullGroup));
processSdf01And57(find(SdfTable.SDF_01, fullGroup), find(SdfTable.SDF_57, fullGroup));
} else if (sdfGroup.get() == SdfGroup.session_Six) {
processSdf04(find(SdfTable.SDF_04, fullGroup));
processSdf13(find(SdfTable.SDF_13, fullGroup));
processSdf08And21(find(SdfTable.SDF_08, fullGroup), find(SdfTable.SDF_21, fullGroup));
processSdf01And57(find(SdfTable.SDF_01, fullGroup), find(SdfTable.SDF_57, fullGroup));
} else { } else {
throw new IllegalStateException("not implemented"); throw new IllegalStateException("not implemented");
} }
@ -109,6 +163,110 @@ public class StatementServiceV2 {
log.info("pair sdf08/sdf21 processed successfully"); log.info("pair sdf08/sdf21 processed successfully");
} }
/**
* fromAccService передается когда пришел ответ от account-service
* в этом случае: по key находим пару в которой сохранен sdf57 запрос и частично выполненный sdf01
* вместо старого sdf01 запроса выполняем новый пришедший от account-service
*/
private void processSdf01And57(StatementRequest sdf01, StatementRequest sdf57) {
Result sdf01Res = processSdf01(sdf01);
removeFirstWithSameTableAndGroupId(sdf01);
if (sdf01Res.getAccountRequests().size() > 0) {
log.info("sdf01 execution wasn't complete, waiting for an answer from account-service");
return;
}
//затем sdf57
processSdf57(sdf57);
removeFirstWithSameTableAndGroupId(sdf57);
reviser.doRevise(sdf01.getGroupId());
//теперь можем продолжить сессию с шага 1
SessionContinueEvent continueSessionBn = new SessionContinueEvent(SdfTable.SDF_01, SdfTable.SDF_57);
kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN_FIRST_PART, continueSessionBn);
log.info("pair sdf01/sdf57 processed successfully");
}
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 = sdf04Executor;
removeFirstWithSameTableAndGroupId(statementRequest);
if (service != null) {
Result res = service.execute(sdfGroup, statementRequest);
SessionContinueEvent continueSessionBn = new SessionContinueEvent(SdfTable.SDF_04);
if (res.isAnyHasError()) {
NotificationNewRequest nRequest = new NotificationNewRequest();
nRequest.setObjectType(ObjectType.rgst.getKey());
nRequest.setPriority(Priority.HIGH.getKey());
nRequest.setComment("Поручения из клиринговой системы имеют неисполненный статус в ответе ДФ-04 из расчетной организации");
kafkaSender.sendRequestToQueue(Consts.NOTIFICATION_NEW, nRequest);
} else { // отправить в SessionMonitor
kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN_FIRST_PART, continueSessionBn);
}
} else {
log.warn("Executor for SDF_04 not set");
}
}
private Result 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()));
} else {
log.debug("accounts were created: {} for generationId {} child generationId {}",
statementRequest.getAccountCreationResults().size(), statementRequest.getGroupId(), statementRequest.getChildGenerationId());
//todo могут ли тут реквесты бегать по кругу, если да, то у нас проблемы
sdfGroup = statementRequest.getAccountCreationResults()
.stream()
// .filter(part -> part.getErrorCode() == null) //fixme эти случае должны попадать в ошибочный sdf02
.map(part -> sdfImdg.getSingleObjectByID(part.getSdfId()))
.collect(Collectors.toList());
}
AbstractExecutor service = sdf01Executor;
Result res = service.execute(sdfGroup, statementRequest);
if (res.getAccountRequests().size() != 0) {
kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF01, StatementService.createAccountsRequest(statementRequest.getGroupId(), res.getAccountRequests(), res.getChildGenerationId()));
return res;
} else if (service.isNeedToSendCommand()) {
service.sendCommand(kafkaSender, res);
return res;
}
return res;
}
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 = sdf57Executor;
Result res = service.execute(sdfGroup, statementRequest);
// finishSendCommand(res, service, 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 = sdf13Executor;
removeFirstWithSameTableAndGroupId(statementRequest);
Result res = service.execute(sdfGroup, statementRequest);
SessionContinueEvent continueSessionBn = new SessionContinueEvent(SdfTable.SDF_13);
if (res.isAnyHasError()) {
NotificationNewRequest nRequest = new NotificationNewRequest();
nRequest.setObjectType(ObjectType.rgst.getKey());
nRequest.setPriority(Priority.HIGH.getKey());
nRequest.setComment("Поручения из клиринговой системы имеют неисполненный статус в ответе ДФ-13 из расчетного депозитария");
kafkaSender.sendRequestToQueue(Consts.NOTIFICATION_NEW, nRequest);
} else {
kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN_FIRST_PART, continueSessionBn);
}
// finishSendCommand(res, service, statementRequest);
}
private void removeFirstWithSameTableAndGroupId(StatementRequest statementRequest) { private void removeFirstWithSameTableAndGroupId(StatementRequest statementRequest) {
Iterator<StatementRequest> iterator = statementRequests.iterator(); Iterator<StatementRequest> iterator = statementRequests.iterator();
Predicate<StatementRequest> eqs = r -> Objects.equals(r.getTable(), statementRequest.getTable()) && Objects.equals(r.getGroupId(), statementRequest.getGroupId()); Predicate<StatementRequest> eqs = r -> Objects.equals(r.getTable(), statementRequest.getTable()) && Objects.equals(r.getGroupId(), statementRequest.getGroupId());
@ -155,7 +313,7 @@ public class StatementServiceV2 {
} }
private Collection<StatementRequest> getFullGroup(StatementRequest stmtReq, SdfGroup group) { private List<StatementRequest> getFullGroup(StatementRequest stmtReq, SdfGroup group) {
List<StatementRequest> reqs = new ArrayList<>(); List<StatementRequest> reqs = new ArrayList<>();
log.debug("trying to find request pair for {}. cashed reqs {}", stmtReq.getTable(), log.debug("trying to find request pair for {}. cashed reqs {}", stmtReq.getTable(),
statementRequests.stream() statementRequests.stream()

View file

@ -5,8 +5,12 @@ import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.regex.Matcher; import java.util.regex.Matcher;
import java.util.regex.Pattern; import java.util.regex.Pattern;
import java.util.stream.Collector;
import java.util.stream.Collectors;
public class TextUtil { public class TextUtil {
public static final Collector<CharSequence, ?, String> join = Collectors.joining(",", "[", "]");
public static boolean isEmpty(String cs) { public static boolean isEmpty(String cs) {
return cs == null || cs.trim().length() == 0; return cs == null || cs.trim().length() == 0;
} }