scheduler-service OUTV шедулер редиректит его в очередь PAYMENT_INSTRUCTION_CLEARING_OUTBOUND_ACTION

This commit is contained in:
AKurakin 2023-08-17 17:49:44 +03:00
parent 25c2be5c52
commit 848c26d69a

View file

@ -3,6 +3,7 @@ package ru.spcex.clearing.scheduler.service;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.InitializingBean;
@ -11,8 +12,10 @@ import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.scheduler.Launcher;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.ActionType;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.payment.PIClearingOutbondActionNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
@ -21,8 +24,11 @@ import ru.spcex.clearing.util.security.UserRoleVerification;
import ru.spcex.platform.enumeration.Task;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.utils.log.ExceptionUtils;
import java.time.Instant;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Future;
import static ru.spcex.platform.utils.enumeration.IEnumKey.getEnumByKey;
@ -74,7 +80,39 @@ public class LauncherService extends QueueConsumer implements InitializingBean {
launcher.setUpdated(created);
launcherMap.insert(launcher);
log.debug("successfully processed, new id {}", launcher.getId());
kafkaProducer.send(new ProducerRecord<>("launcher-" + launcher.getTask(), userRequest));
if (Task.dbfExport_OUTV.equalsByKey(req.getTaskName())) {
processRequestOUTV(Task.dbfExport_OUTV, req);
} else {
kafkaProducer.send(new ProducerRecord<>("launcher-" + launcher.getTask(), userRequest));
}
}
private void processRequestOUTV(Task dbfExport__OUTV, LauncherCommandRequest launcherNew) {
log.debug("Redirect {} to queue {}", dbfExport__OUTV, Consts.PAYMENT_INSTRUCTION_CLEARING_OUTBOUND_ACTION);
var req = new PIClearingOutbondActionNewRequest();
req.setSenderId(launcherNew.getSenderId());
req.setAddresseeId(launcherNew.getAddresseeId());
req.setPaymentPurpose(launcherNew.getPaymentPurpose());
req.setCreditLeg_amount(launcherNew.getCreditLeg_amount());
req.setCreditLeg_accountId(launcherNew.getCreditLeg_accountId());
req.setDebitLeg_accountId(launcherNew.getDebitLeg_accountId());
final String destination = Consts.PAYMENT_INSTRUCTION_CLEARING_OUTBOUND_ACTION;
BaseRequest<Object> request = new BaseRequest<>();
request.setId(null);
request.setActionType(ActionType.NEW);
request.setRequestPayload(req);
request.setUserId(launcherNew.getUserId());
Future<RecordMetadata> send = kafkaProducer.send(new ProducerRecord<>(destination, request));
try {
send.get();
log.debug("Request id={} by user {} send to queue \"{}\". Action: {}", request.getId(), request.getUserId(), destination, request.getActionType());
} catch (InterruptedException e) {
log.error("Request id={} by user {} send to queue \"{}\". Action: {}; error send: thread interrupted", request.getId(), request.getUserId(), destination, request.getActionType());
Thread.currentThread().interrupt();
} catch (ExecutionException e) {
log.error("Request id={} by user {} send to queue \"{}\". Action: {}; error send: {}", request.getId(), request.getUserId(), destination, request.getActionType(), ExceptionUtils.getStackTrace(e.getCause()));
}
}
}