http://jira.mfd.msk:8088/browse/CLS-130 добавил обновление PaymentInstruction.transactionStatus по sdf04

This commit is contained in:
ialbert 2022-12-06 15:24:30 +03:00
parent 0b24bae280
commit abf51317ca
10 changed files with 212 additions and 102 deletions

View file

@ -6,118 +6,33 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
import ru.clearing.classes.statics.data.sdf.SDf03;
import ru.clearing.classes.statics.data.sdf.SDf11;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SdfClearingRequest;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
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 java.util.AbstractMap;
import java.util.Comparator;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.stream.Collectors;
@Service
@EnableScheduling
public class ClearingService {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<PaymentInstruction> paymentImdgs;
private final ExecutorService executor;
private final ImdgId idGenerator;
private final PaymentInstructionSorter senderGroupSorter;
private final Imdg<SDf03> sdf03Imdg;
private final Imdg<SDf11> sdf11Imdg;
private final KafkaSender kafkaSender;
private final SdfCreatorBySTLDPayment sdfCreator;
private final PaymentUpdateBySdf04 paymentUpdater;
@Autowired
public ClearingService(ImdgProvider imdgProvider, PaymentInstructionSorter senderGroupSorter, KafkaSender kafkaSender) {
this.paymentImdgs = imdgProvider.getImdg(IMDGDistributedNames.Map_PaymentInstruction, PaymentInstruction.class);
this.sdf03Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf03, SDf03.class);
this.sdf11Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf11, SDf11.class);
this.idGenerator = imdgProvider.getImdgIdGenerator();
this.senderGroupSorter = senderGroupSorter;
this.kafkaSender = kafkaSender;
public ClearingService(SdfCreatorBySTLDPayment sdfCreator, PaymentUpdateBySdf04 paymentUpdater) {
this.sdfCreator = sdfCreator;
this.paymentUpdater = paymentUpdater;
this.executor = Executors.newSingleThreadExecutor();
}
@Scheduled(cron = "${clearing-service.scheduler.check-payment-instruction}")
public void run() {
executor.execute(this::createSdfFromPaymentInstructionSTLD);
public void sdfCreate() {
log.info("creating sdf03/11 from STLD payments task added to queue");
executor.execute(sdfCreator::createSdfFromPaymentInstructionSTLD);
}
private void createSdfFromPaymentInstructionSTLD() {
final boolean[] anyError = {false};
//generationId для созадаваемых Sdf03/Sdf11
Long generationId = idGenerator.nextId();
//выгружаем PaymentInstructions с нужным статусом
Map<Long, PaymentBatchInfo> paymentBySender = paymentImdgs.getCollectionObjectsByFieldValues(
Map.of("transactionStatus", TransactionStatus.stld.getKey()))
.stream()
//группируем по компаниям (fixme sorted убрать?)
.sorted(Comparator.comparing(PaymentInstruction::getSenderId))
.collect(Collectors.groupingBy(PaymentInstruction::getSenderId))
.entrySet()
.stream()
//результатом работы senderGroupSorter будет Map<senderId -> PaymentBatchInfo>
//PaymentBatchInfo содержит возможную ошибку, при необходимости отсортированные Payment
//тип ClearingMemberCategory
.map(entry -> {
Long senderId = entry.getKey();
List<PaymentInstruction> pmtInstrcs = entry.getValue();
PaymentBatchInfo senderInfo = senderGroupSorter.sortCompanyPayments(generationId, senderId, pmtInstrcs);
if (senderInfo.getError() != null) {
anyError[0] = true;
}
return new AbstractMap.SimpleEntry<>(senderId, senderInfo);
})
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue));
//save SDF03/SDF11
for (var entry : paymentBySender.entrySet()) {
PaymentBatchInfo senderPayments = entry.getValue();
saveSdfAnSendToKafka(senderPayments, generationId);
}
//update PaymentInstruction.transactionStatus
for (var entry : paymentBySender.entrySet()) {
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 -> {
paymentInstruction.setTransactionStatus(stat.getKey());
paymentImdgs.update(paymentInstruction);
});
}
}
}
private void saveSdfAnSendToKafka(PaymentBatchInfo batch, Long generationId) {
SdfClearingRequest kafkaMessage = new SdfClearingRequest();
kafkaMessage.setGroupId(generationId);
switch (batch.getCategoryD()) {
case I -> {
batch.getOrderedPaymentInstructions()
.map(paymentInstruction -> Sdf03Builder.buildSdf03(paymentInstruction, generationId))
.forEach(sdf03Imdg::insert);
kafkaSender.sendRequestToQueue(Consts.SDF03_PROCESS, kafkaMessage);
}
case B -> {
batch.getOrderedPaymentInstructions()
.map(paymentInstruction -> Sdf11Builder.buildSdf11(paymentInstruction, generationId))
.forEach(sdf11Imdg::insert);
kafkaSender.sendRequestToQueue(Consts.SDF11_PROCESS, kafkaMessage);
}
}
public void paymentUpdateBySdf04(Long sdf04GroupId) {
log.info("adding sdf03/11 from STLD payments task to queue");
executor.execute(() -> paymentUpdater.updatePayments(sdf04GroupId));
}
}

View file

@ -0,0 +1,47 @@
package ru.spcex.clearing.service;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
import ru.clearing.classes.statics.data.sdf.SDf04;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.enumeration.TransactionStatus;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import java.util.Collection;
import java.util.Map;
@Component
public class PaymentUpdateBySdf04 {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<PaymentInstruction> paymentImdgs;
private final Imdg<SDf04> sdf04Imdg;
public PaymentUpdateBySdf04(ImdgProvider imdgProvider) {
this.paymentImdgs = imdgProvider.getImdg(IMDGDistributedNames.Map_PaymentInstruction, PaymentInstruction.class);
this.sdf04Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf04, SDf04.class);
}
public void updatePayments(Long sdf04GroupId) {
Collection<SDf04> sDf04s = sdf04Imdg.getCollectionObjectsByFieldValues(Map.of("generationId", sdf04GroupId));
for (SDf04 sDf04 : sDf04s) {
String docnm_ref = sDf04.getDocnm_ref();
long paymentId;
try {
paymentId = Long.parseLong(docnm_ref);
} catch (NumberFormatException e) {
log.warn("sdf04.id={} cannot parse docnm_ref={} as payment id", sDf04.getId(), sDf04.getDocnm_ref());
continue;
}
PaymentInstruction payment = paymentImdgs.getSingleObjectByFieldValues(Map.of("id", paymentId));
if ("OK!".equals(sDf04.getImp_result())) {
payment.setTransactionStatus(TransactionStatus.ok.getKey());
} else {
payment.setTransactionStatus(TransactionStatus.fail.getKey());
}
paymentImdgs.update(payment);
}
}
}

View file

@ -0,0 +1,28 @@
package ru.spcex.clearing.service;
import org.apache.kafka.clients.consumer.Consumer;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.stereotype.Service;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.Sdf04Request;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
@Service
public class Sdf04Receiver extends QueueConsumer implements InitializingBean {
private final ClearingService clearingService;
public Sdf04Receiver(Consumer<String, Object> kafkaQueue, ClearingService clearingService) {
super(kafkaQueue);
this.clearingService = clearingService;
}
@Override
public void afterPropertiesSet() {
callback(Sdf04Request.class)
.setConsumer(event -> {
Sdf04Request requestPayload = event.getRequestPayload();
clearingService.paymentUpdateBySdf04(requestPayload.getGroupId());
})
.forDestination(Consts.SDF04_PROCESS, callbacks::put);
init();
}
}

View file

@ -0,0 +1,112 @@
package ru.spcex.clearing.service;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
import ru.clearing.classes.statics.data.sdf.SDf03;
import ru.clearing.classes.statics.data.sdf.SDf11;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SdfClearingRequest;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.order.PaymentBatchInfo;
import ru.spcex.clearing.service.order.PaymentInstructionSorter;
import ru.spcex.clearing.service.order.Sdf03Builder;
import ru.spcex.clearing.service.order.Sdf11Builder;
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 java.util.AbstractMap;
import java.util.Comparator;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;
@Component
public class SdfCreatorBySTLDPayment {
private final Imdg<PaymentInstruction> paymentImdgs;
private final ImdgId idGenerator;
private final PaymentInstructionSorter senderGroupSorter;
private final Imdg<SDf03> sdf03Imdg;
private final Imdg<SDf11> sdf11Imdg;
private final KafkaSender kafkaSender;
@Autowired
public SdfCreatorBySTLDPayment(ImdgProvider imdgProvider, PaymentInstructionSorter senderGroupSorter, KafkaSender kafkaSender) {
this.paymentImdgs = imdgProvider.getImdg(IMDGDistributedNames.Map_PaymentInstruction, PaymentInstruction.class);
this.sdf03Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf03, SDf03.class);
this.sdf11Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf11, SDf11.class);
this.idGenerator = imdgProvider.getImdgIdGenerator();
this.senderGroupSorter = senderGroupSorter;
this.kafkaSender = kafkaSender;
}
public void createSdfFromPaymentInstructionSTLD() {
final boolean[] anyError = {false};
//generationId для созадаваемых Sdf03/Sdf11
Long generationId = idGenerator.nextId();
//выгружаем PaymentInstructions с нужным статусом
Map<Long, PaymentBatchInfo> paymentBySender = paymentImdgs.getCollectionObjectsByFieldValues(
Map.of("transactionStatus", TransactionStatus.stld.getKey()))
.stream()
//группируем по компаниям (fixme sorted убрать?)
.sorted(Comparator.comparing(PaymentInstruction::getSenderId))
.collect(Collectors.groupingBy(PaymentInstruction::getSenderId))
.entrySet()
.stream()
//результатом работы senderGroupSorter будет Map<senderId -> PaymentBatchInfo>
//PaymentBatchInfo содержит возможную ошибку, при необходимости отсортированные Payment
//тип ClearingMemberCategory
.map(entry -> {
Long senderId = entry.getKey();
List<PaymentInstruction> pmtInstrcs = entry.getValue();
PaymentBatchInfo senderInfo = senderGroupSorter.sortCompanyPayments(generationId, senderId, pmtInstrcs);
if (senderInfo.getError() != null) {
anyError[0] = true;
}
return new AbstractMap.SimpleEntry<>(senderId, senderInfo);
})
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue));
//save SDF03/SDF11
for (var entry : paymentBySender.entrySet()) {
PaymentBatchInfo senderPayments = entry.getValue();
saveSdfAnSendToKafka(senderPayments, generationId);
}
//update PaymentInstruction.transactionStatus
for (var entry : paymentBySender.entrySet()) {
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 -> {
paymentInstruction.setTransactionStatus(stat.getKey());
paymentImdgs.update(paymentInstruction);
});
}
}
}
private void saveSdfAnSendToKafka(PaymentBatchInfo batch, Long generationId) {
SdfClearingRequest kafkaMessage = new SdfClearingRequest();
kafkaMessage.setGroupId(generationId);
switch (batch.getCategoryD()) {
case I -> {
batch.getOrderedPaymentInstructions()
.map(paymentInstruction -> Sdf03Builder.buildSdf03(paymentInstruction, generationId))
.forEach(sdf03Imdg::insert);
kafkaSender.sendRequestToQueue(Consts.SDF03_PROCESS, kafkaMessage);
}
case B -> {
batch.getOrderedPaymentInstructions()
.map(paymentInstruction -> Sdf11Builder.buildSdf11(paymentInstruction, generationId))
.forEach(sdf11Imdg::insert);
kafkaSender.sendRequestToQueue(Consts.SDF11_PROCESS, kafkaMessage);
}
}
}
}

View file

@ -1,4 +1,4 @@
package ru.spcex.clearing.service;
package ru.spcex.clearing.service.order;
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
import ru.spcex.platform.enumeration.ClearingMemberCategoryD;

View file

@ -1,4 +1,4 @@
package ru.spcex.clearing.service;
package ru.spcex.clearing.service.order;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@ -40,7 +40,7 @@ public class PaymentInstructionSorter {
this.messageResolver = messageResolver;
}
PaymentBatchInfo sortCompanyPayments(Long generationId, Long senderId, List<PaymentInstruction> payments) {
public PaymentBatchInfo sortCompanyPayments(Long generationId, Long senderId, List<PaymentInstruction> payments) {
PaymentBatchInfo batchInfo = new PaymentBatchInfo();
log.info("processing PaymentInstruction's generationId={} senderId={} size={}", generationId, senderId, payments.size());
ClearingMemberCategory category = clrngMmbrImdg.getSingleObjectByFieldValues(Map.of("companyId", senderId));

View file

@ -1,7 +1,8 @@
package ru.spcex.clearing.service;
package ru.spcex.clearing.service.order;
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
import ru.clearing.classes.statics.data.sdf.SDf03;
import ru.spcex.clearing.service.SpecifUtil;
import ru.spcex.platform.enumeration.Sender;
import ru.spcex.platform.utils.time.TimeUtil;

View file

@ -1,7 +1,8 @@
package ru.spcex.clearing.service;
package ru.spcex.clearing.service.order;
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
import ru.clearing.classes.statics.data.sdf.SDf11;
import ru.spcex.clearing.service.SpecifUtil;
import ru.spcex.platform.enumeration.Sender;
import ru.spcex.platform.utils.time.TimeUtil;

View file

@ -3,7 +3,7 @@ package ru.spcex.platform.enumeration;
import ru.spcex.platform.utils.enumeration.IEnumKey;
public enum TransactionStatus implements IEnumKey {
stld("STLD"), notSent("NSNT"), cher("CHER"), sent("SENT");
stld("STLD"), notSent("NSNT"), cher("CHER"), sent("SENT"), ok("OK"), fail("FAIL");
private final String key;

View file

@ -50,6 +50,12 @@ public class QueueConsumer implements AutoCloseable {
this.supportStartOffsetTimeWindow = false;
}
/**
* используем этот конструктор, если хотим класть
* в кафку "ответ" - информацию о статусе обработки команд
* @param kafkaQueue
* @param kafkaResponseQueue
*/
public QueueConsumer(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaResponseQueue) {
this(kafkaQueue);
this.producer = kafkaResponseQueue;