ialbert 2023-02-06 20:06:34 +03:00
parent 0e8ea7b595
commit 09605b610d

View file

@ -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<PaymentInstruction> paymentImdgs;
private final Imdg<ClearingMemberCategory> clearingMemberCategoryImdg;
private final ImdgId idGenerator;
private final PaymentInstructionSorter senderGroupSorter;
private final Imdg<SDf03> 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<PaymentInstruction> paymentInstructionFound = paymentImdgs.getCollectionObjectsByFieldValues(
Map.of("transactionStatus", TransactionStatus.stld.getKey()));
log.info("Found {} PaymentInstruction by status {}", paymentInstructionFound.size(), TransactionStatus.stld.getKey());
List<AbstractMap.SimpleEntry<Long, PaymentBatchInfo>> paymentBySenderLst = new ArrayList<>();
paymentInstructionFound.stream()
//группируем по компаниям
.sorted(Comparator.comparing(PaymentInstruction::getSenderId))
.collect(Collectors.groupingBy(PaymentInstruction::getSenderId))
/*Map<GroupWrapper, List<PaymentInstruction>> byOrder = */
List<Map.Entry<GroupWrapper, Collection<List<PaymentInstruction>>>> 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<senderId -> PaymentBatchInfo>
//PaymentBatchInfo содержит возможную ошибку, при необходимости отсортированные Payment
//тип ClearingMemberCategory
.forEachOrdered(entry -> {
Long senderId = entry.getKey();
List<PaymentInstruction> pmtInstrcs = entry.getValue();
List<PaymentBatchInfo> 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<group, list<paymentInstruction>>, разбиваем список на группы по senderId
.map((Function<Map.Entry<GroupWrapper, List<PaymentInstruction>>, Map.Entry<GroupWrapper, Collection<List<PaymentInstruction>>>>) entry -> {
Collection<List<PaymentInstruction>> 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<List<PaymentInstruction>> bySenderGroups = entry.getValue();
for (List<PaymentInstruction> 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<PaymentInstruction> 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<GroupWrapper> {
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());
}
}