From 21df705094e889a6500a19dda969dd79a8209c13 Mon Sep 17 00:00:00 2001 From: ialbert Date: Tue, 5 Sep 2023 19:14:17 +0300 Subject: [PATCH] =?UTF-8?q?SDF=20=D0=B3=D1=80=D1=83=D0=BF=D0=BF=D1=8B=20?= =?UTF-8?q?=D0=B4=D0=BB=D1=8F=20=D1=81=D0=B5=D1=81=D1=81=D0=B8=D0=B9=20?= =?UTF-8?q?=D0=B8=D0=B7=D0=B1=D0=B0=D0=B2=D0=B8=D0=BB=D1=81=D1=8F=20=D0=BE?= =?UTF-8?q?=D1=82=20StatementService.java?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../clearing/service/EventsReceiver.java | 8 +- .../clearing/service/StatementService.java | 28 ++- .../ru/spcex/clearing/statement/SdfGroup.java | 7 +- .../clearing/statement/SdfGroupManager.java | 61 ++++++ .../statement/StatementServiceV2.java | 174 +++++++++++++++++- .../spcex/platform/utils/text/TextUtil.java | 4 + 6 files changed, 266 insertions(+), 16 deletions(-) create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/SdfGroupManager.java diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java index a4d10dff0..db3b9959c 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java @@ -29,6 +29,7 @@ import ru.spcex.clearing.service.executors.Sdf10Executor; import ru.spcex.clearing.service.payment.PaymentInstructionOutboundService; import ru.spcex.clearing.session.stage.*; import ru.spcex.clearing.session.stage.impl.BalanceRevise; +import ru.spcex.clearing.statement.StatementServiceV2; import ru.spcex.platform.enumeration.Task; import ru.spcex.platform.utils.enumeration.IMessageResolver; import ru.spcex.platform.utils.error.ValidationException; @@ -53,6 +54,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { private final Sdf10Executor sdf10Executor; private final BalanceRevise balanceRevise; private final Sdf05Sender sdf05Sender; + private final StatementServiceV2 statementService; private final PaymentInstructionOutboundService pmtOutboundService; @Autowired @@ -64,7 +66,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { SecondaryAuctionT0Session secondaryAuctionT0Session, PrimaryAuctionB0Session primaryAuctionB0Session, PrimaryAuctionT0Session primaryAuctionT0Session, IntermediateMkrSession intermediateMkrSession, FinalMkrSession finalMkrSession, ReturnDepositSession returnDepositSession, SessionManager sessionManager, Sdf06Executor sdf06Executor, - Sdf10Executor sdf10Executor, BalanceRevise balanceRevise, Sdf05Sender sdf05Sender, PaymentInstructionOutboundService pmtOutboundService) { + Sdf10Executor sdf10Executor, BalanceRevise balanceRevise, Sdf05Sender sdf05Sender, StatementServiceV2 statementService, PaymentInstructionOutboundService pmtOutboundService) { super(kafkaQueue, kafkaResponseQueue); this.errorResolver = errorResolver; this.clearingService = clearingService; @@ -81,6 +83,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { this.sdf10Executor = sdf10Executor; this.balanceRevise = balanceRevise; this.sdf05Sender = sdf05Sender; + this.statementService = statementService; this.pmtOutboundService = pmtOutboundService; } @@ -116,6 +119,9 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { callback(PIClearingOutbondActionNewRequest.class) .setFunction(pmtOutboundService::sendOutBoundPayment) .forDestination(Consts.PAYMENT_INSTRUCTION_CLEARING_OUTBOUND_ACTION, callbacks::put); + callback(StatementRequest.class) + .setConsumer(statementService::processReq) + .forDestination(Consts.STATEMENT_PROCESS, callbacks::put); callback(SessionContinueEvent.class) .setConsumer(req -> { primaryAuctionBnSession.continueSession(req); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java index 381f7bbcd..8efaa429a 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java @@ -4,9 +4,7 @@ import org.apache.kafka.clients.consumer.Consumer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; -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.*; import ru.spcex.clearing.imdg.IMDGDistributedNames; 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.stream.Collectors; -@Service +//@Service +@Deprecated public class StatementService extends QueueConsumer implements InitializingBean { private final Logger log = LoggerFactory.getLogger(getClass()); @@ -49,7 +48,7 @@ public class StatementService extends QueueConsumer implements InitializingBean */ private final Map> pairOfSdfRequest = new HashMap<>(); - @Autowired +// @Autowired public StatementService(Consumer kafkaQueue, ImdgProvider imdgProvider, KafkaSender kafkaSender, @@ -75,11 +74,28 @@ public class StatementService extends QueueConsumer implements InitializingBean .setConsumer(systemRequest -> { StatementRequest payload = systemRequest.getRequestPayload(); 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()); 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); } else if (payload.isContinueSdf()) { processAccountAnswer(systemRequest); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/SdfGroup.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/SdfGroup.java index f1374118d..d48c42cee 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/SdfGroup.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/SdfGroup.java @@ -13,7 +13,12 @@ public enum SdfGroup { //------ session groups ------ 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 group; diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/SdfGroupManager.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/SdfGroupManager.java new file mode 100644 index 000000000..877681bc8 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/SdfGroupManager.java @@ -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 sessImdg; + + @Autowired + public SdfGroupManager(ImdgProvider imdgProvider) { + sessImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class); + } + + public Optional getGroup(SdfTable table) { + ImdgPredicateBuilder pb = sessImdg.predicateBuilder(); + ImdgPredicate prdct = pb.and( + pb.equals("workflowStatus", SessionStatus.ACTV.getKey()), + pb.equals("sessionStatus", TaskType.FormingPaymentInstruction.getKey()) + ); + Collection 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()); + } + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/StatementServiceV2.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/StatementServiceV2.java index bd5011878..2814e108d 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/StatementServiceV2.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/StatementServiceV2.java @@ -3,20 +3,25 @@ package ru.spcex.clearing.statement; 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.clearing.classes.statics.data.sdf.*; 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.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.service.StatementService; import ru.spcex.clearing.service.executors.*; 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.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.utils.text.TextUtil; import java.util.*; import java.util.function.Predicate; @@ -33,6 +38,12 @@ public class StatementServiceV2 { private final Sdf10Executor sdf10Executor; private final Sdf55Executor sdf55Executor; 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, KafkaSender kafkaSender, @@ -40,7 +51,7 @@ public class StatementServiceV2 { Sdf21Executor sdf21Executor, Sdf10Executor sdf10Executor, Sdf55Executor sdf55Executor, - Sdf20Executor sdf20Executor) { + Sdf20Executor sdf20Executor, Sdf01Executor sdf01Executor, Sdf57Executor sdf57Executor, Sdf04Executor sdf04Executor, Sdf13Executor sdf13Executor, Reviser reviser, SdfGroupManager grpMng) { this.imdgProvider = imdgProvider; this.kafkaSender = kafkaSender; this.sdf08Executor = sdf08Executor; @@ -48,35 +59,78 @@ public class StatementServiceV2 { this.sdf10Executor = sdf10Executor; this.sdf55Executor = sdf55Executor; 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 systemRequest) { SdfTable table = systemRequest.getRequestPayload().getTable(); 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); - Optional sdfGroup = SdfGroup.groupByTable(table); + Optional sdfGroup = grpMng.getGroup(table); if (sdfGroup.isEmpty()) { log.debug("table {} is not paired with any other SDF", table); //здесь switch case одиночные методы switch (table) { + case SDF_04 -> processSdf04(systemRequest.getRequestPayload()); 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_55 -> sdf55Executor.execute(systemRequest); default -> log.error("unknown table {}", table); } return; } - Collection 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 fullGroup = getFullGroup(systemRequest.getRequestPayload(), sdfGroup.get()); statementRequests.add(systemRequest.getRequestPayload()); if (fullGroup.isEmpty()) { 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; } log.debug("find full set of SDF requests: {}", fullGroup.stream() .map(stReq -> stReq.getTable().getKey() + " generationId=" + stReq.getGroupId()) - .collect(Collectors.joining(",", "[", "]"))); + .collect(TextUtil.join)); if (sdfGroup.get() == SdfGroup.Sdf08And21) { 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 { throw new IllegalStateException("not implemented"); } @@ -109,6 +163,110 @@ public class StatementServiceV2 { 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 sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf04, SDf04.class); + Collection 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 sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf01, SDf01.class); + Collection 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 sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf57, SDf57.class); + Collection 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 sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf13, SDf13.class); + Collection 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) { Iterator iterator = statementRequests.iterator(); Predicate eqs = r -> Objects.equals(r.getTable(), statementRequest.getTable()) && Objects.equals(r.getGroupId(), statementRequest.getGroupId()); @@ -155,7 +313,7 @@ public class StatementServiceV2 { } - private Collection getFullGroup(StatementRequest stmtReq, SdfGroup group) { + private List getFullGroup(StatementRequest stmtReq, SdfGroup group) { List reqs = new ArrayList<>(); log.debug("trying to find request pair for {}. cashed reqs {}", stmtReq.getTable(), statementRequests.stream() diff --git a/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/text/TextUtil.java b/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/text/TextUtil.java index 3d43ea1b5..7c1aa5598 100644 --- a/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/text/TextUtil.java +++ b/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/text/TextUtil.java @@ -5,8 +5,12 @@ import java.util.List; import java.util.Map; import java.util.regex.Matcher; import java.util.regex.Pattern; +import java.util.stream.Collector; +import java.util.stream.Collectors; public class TextUtil { + public static final Collector join = Collectors.joining(",", "[", "]"); + public static boolean isEmpty(String cs) { return cs == null || cs.trim().length() == 0; }