diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/SdfCreatorBySTLDPayment.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/SdfCreatorBySTLDPayment.java index cc1b8a607..89b515660 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/SdfCreatorBySTLDPayment.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/SdfCreatorBySTLDPayment.java @@ -4,6 +4,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; +import ru.clearing.classes.statics.data.generated.ClearingMemberCategory; import ru.clearing.classes.statics.data.payment.PaymentInstruction; import ru.clearing.classes.statics.data.sdf.SDf03; import ru.clearing.classes.statics.data.sdf.SDf11; @@ -13,20 +14,25 @@ import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SdfClearingReque import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.service.builder.Sdf03Builder; import ru.spcex.clearing.service.builder.Sdf11Builder; -import ru.spcex.clearing.service.order.PaymentBatchInfo; import ru.spcex.clearing.service.order.PaymentInstructionSorter; +import ru.spcex.platform.enumeration.ClearingCategory; import ru.spcex.platform.enumeration.TransactionStatus; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgId; import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.utils.collection.Pair; +import ru.spcex.platform.utils.enumeration.IEnumKey; import java.util.*; +import java.util.function.Function; import java.util.stream.Collectors; +import java.util.stream.Stream; @Component public class SdfCreatorBySTLDPayment { private Logger log = LoggerFactory.getLogger(getClass()); private final Imdg paymentImdgs; + private final Imdg clearingMemberCategoryImdg; private final ImdgId idGenerator; private final PaymentInstructionSorter senderGroupSorter; private final Imdg sdf03Imdg; @@ -36,6 +42,7 @@ public class SdfCreatorBySTLDPayment { @Autowired public SdfCreatorBySTLDPayment(ImdgProvider imdgProvider, PaymentInstructionSorter senderGroupSorter, KafkaSender kafkaSender) { this.paymentImdgs = imdgProvider.getImdg(IMDGDistributedNames.Map_PaymentInstruction, PaymentInstruction.class); + this.clearingMemberCategoryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingMemberCategory, ClearingMemberCategory.class); this.sdf03Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf03, SDf03.class); this.sdf11Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf11, SDf11.class); this.idGenerator = imdgProvider.getImdgIdGenerator(); @@ -43,76 +50,134 @@ public class SdfCreatorBySTLDPayment { this.kafkaSender = kafkaSender; } + public void createSdfFromPaymentInstructionSTLD() { - final boolean[] anyError = {false}; //generationId для созадаваемых Sdf03/Sdf11 Long generationId = idGenerator.nextId(); //выгружаем PaymentInstructions с нужным статусом Collection paymentInstructionFound = paymentImdgs.getCollectionObjectsByFieldValues( Map.of("transactionStatus", TransactionStatus.stld.getKey())); log.info("Found {} PaymentInstruction by status {}", paymentInstructionFound.size(), TransactionStatus.stld.getKey()); - List> paymentBySenderLst = new ArrayList<>(); - paymentInstructionFound.stream() - //группируем по компаниям - .sorted(Comparator.comparing(PaymentInstruction::getSenderId)) - .collect(Collectors.groupingBy(PaymentInstruction::getSenderId)) + /*Map> byOrder = */ + List>>> allEntries = paymentInstructionFound.stream().flatMap(paymentInstruction -> { + ClearingCategory clearingCategory = getClearingCategory(paymentInstruction); + Long senderId = paymentInstruction.getSenderId(); + if (senderId == null) { + return Stream.empty(); + } else { + return Stream.of(new Pair<>(new GroupWrapper(clearingCategory, senderId), paymentInstruction)); + } + //группируем по equals/hashCode группы + }).collect(Collectors.groupingBy(Pair::getFirst, Collectors.mapping(Pair::getSecond, Collectors.toList()))) .entrySet() .stream() - //результатом работы senderGroupSorter будет Map PaymentBatchInfo> - //PaymentBatchInfo содержит возможную ошибку, при необходимости отсортированные Payment - //тип ClearingMemberCategory - .forEachOrdered(entry -> { - Long senderId = entry.getKey(); - List pmtInstrcs = entry.getValue(); - List senderInfos = senderGroupSorter.sortCompanyPayments(generationId, senderId, pmtInstrcs); - for (PaymentBatchInfo senderInfo:senderInfos) { - if (senderInfo.getError() != null) { - anyError[0] = true; - } - paymentBySenderLst.add(new AbstractMap.SimpleEntry<>(senderId, senderInfo)); - } - }); - //save SDF03/SDF11 - log.debug("Sending save {} sdf03/sdf11", paymentBySenderLst.size()); - for (var entry : paymentBySenderLst) { - PaymentBatchInfo senderPayments = entry.getValue(); - saveSdfAnSendToKafka(senderPayments, generationId); - } - //update PaymentInstruction.transactionStatus - for (var entry : paymentBySenderLst) { - PaymentBatchInfo batch = entry.getValue(); - //все PaymentInstruction.transactionStatus в batch с error != null - //уже проапдейтились в методе sortCompanyPayments - if (batch.getError() == null) { - TransactionStatus stat = anyError[0] ? TransactionStatus.notSent : TransactionStatus.sent; - batch.getOrderedPaymentInstructions() - .forEach(paymentInstruction -> { + //сортируем по группе compareTo + .sorted(Map.Entry.comparingByKey()) + //преобразуем поток entry>, разбиваем список на группы по senderId + .map((Function>, Map.Entry>>>) entry -> { + Collection> values = entry.getValue().stream().collect(Collectors.groupingBy(PaymentInstruction::getSenderId)).values(); + return new AbstractMap.SimpleEntry<>(entry.getKey(), values); + }).toList(); + + allEntries.forEach(entry -> { + log.info("Group category {}, index {}", entry.getKey().category, entry.getKey().getSortIndex()); + GroupWrapper group = entry.getKey(); + Collection> bySenderGroups = entry.getValue(); + for (List bySenderGroup : bySenderGroups) { + Long senderIdSpecificToGroup = bySenderGroup.iterator().next().getSenderId(); + log.info("Group senderId {}", senderIdSpecificToGroup); + saveSdfAnSendToKafka(senderIdSpecificToGroup, group.category, bySenderGroup, generationId); + } + }); + allEntries.forEach(entry -> { + //category group + GroupWrapper group = entry.getKey(); + log.info("updating transactionStatus Group category {} index {}", + group.category, group.getSortIndex()); + TransactionStatus stat; + if (group.category.equals(ClearingCategory.I) + || group.category.equals(ClearingCategory.V) + || group.category.equals(ClearingCategory.B)) { + stat = TransactionStatus.sent; + } else { + stat = TransactionStatus.notSent; + } + entry.getValue() + .forEach(senderInstructions -> { + senderInstructions.forEach(paymentInstruction -> { paymentInstruction.setTransactionStatus(stat.getKey()); paymentImdgs.update(paymentInstruction); }); - } - } - + }); + }); } - private void saveSdfAnSendToKafka(PaymentBatchInfo batch, Long generationId) { + private void saveSdfAnSendToKafka(Long senderId, ClearingCategory category, List batch, Long generationId) { SdfClearingRequest kafkaMessage = new SdfClearingRequest(); kafkaMessage.setGroupId(generationId); Long requestId = null; - switch (batch.getCategoryD()) { - case I -> { - batch.getOrderedPaymentInstructions() - .map(paymentInstruction -> Sdf03Builder.buildSdf03(paymentInstruction, generationId)) - .forEach(sdf03Imdg::insert); - requestId = kafkaSender.sendRequestToQueue(Consts.SDF03_PROCESS, kafkaMessage); - } - case B -> { - batch.getOrderedPaymentInstructions() - .map(paymentInstruction -> Sdf11Builder.buildSdf11(paymentInstruction, generationId)) - .forEach(sdf11Imdg::insert); - requestId = kafkaSender.sendRequestToQueue(Consts.SDF11_PROCESS, kafkaMessage); - } + if (category.equals(ClearingCategory.I) || category.equals(ClearingCategory.V)) { + batch.stream() + .map(paymentInstruction -> Sdf03Builder.buildSdf03(paymentInstruction, generationId)) + .forEach(sdf03Imdg::insert); + requestId = kafkaSender.sendRequestToQueue(Consts.SDF03_PROCESS, kafkaMessage); + } else if (category.equals(ClearingCategory.B)) { + batch.stream() + .map(paymentInstruction -> Sdf11Builder.buildSdf11(paymentInstruction, generationId)) + .forEach(sdf11Imdg::insert); + requestId = kafkaSender.sendRequestToQueue(Consts.SDF11_PROCESS, kafkaMessage); + } else { + log.error("Unknown action (sdf03/sdf11) senderId {}, generationId {}, category {}", senderId, generationId, category); } - log.debug("Send to kafka command, CategoryD={}, requestId={}", batch.getCategoryD(), requestId); + log.debug("Send to kafka command, CategoryD={}, requestId={}", category, requestId); + } + + private static class GroupWrapper implements Comparable { + private Long senderId; + private ClearingCategory category; + + public GroupWrapper(ClearingCategory clearingCategory, Long senderId) { + this.category = clearingCategory; + this.senderId = senderId; + } + + private Integer getSortIndex() { + if (category == ClearingCategory.I) + return 0; + if (senderId == 1L) + return 1; + if (category == ClearingCategory.V) + return 2; + if (category == ClearingCategory.B) + return 3; + return 4; + } + + @Override + public int compareTo(GroupWrapper o) { + Integer thisIndex = getSortIndex(); + Integer otherIndex = o.getSortIndex(); + return thisIndex.compareTo(otherIndex); + } + + @Override + public boolean equals(Object o) { + if (this == o) return true; + if (o == null || getClass() != o.getClass()) return false; + GroupWrapper that = (GroupWrapper) o; + return Objects.equals(getSortIndex(), that.getSortIndex()); + } + + @Override + public int hashCode() { + return getSortIndex(); + } + } + + private ClearingCategory getClearingCategory(PaymentInstruction paymentInstruction) { + ClearingMemberCategory category = clearingMemberCategoryImdg.getSingleObjectByFieldValues( + Map.of("companyId", paymentInstruction.getSenderId())); + return IEnumKey.getEnumByKeyOrUndefined(ClearingCategory.class, + category.getClearingMemberCategory()); } }