This commit is contained in:
parent
0d35f74b1c
commit
ed47b5291b
11 changed files with 380 additions and 177 deletions
|
|
@ -46,4 +46,16 @@ public class KafkaConfig {
|
||||||
})
|
})
|
||||||
.build();
|
.build();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
@Bean("kafkaSenderWithoutRequestInfo")
|
||||||
|
public KafkaSender kafkaSenderWithoutRequestInfo(Producer<String, Object> kafkaProducer, ImdgProvider imdgProvider) {
|
||||||
|
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
|
||||||
|
return KafkaSender
|
||||||
|
.setup()
|
||||||
|
.producer(kafkaProducer)
|
||||||
|
.idGenerator(imdgIdGenerator::nextId)
|
||||||
|
.saveRequestInfo(false)
|
||||||
|
.build();
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -7,9 +7,13 @@ import ru.clearing.classes.statics.data.account.AccountBalance;
|
||||||
import ru.clearing.classes.statics.data.company.Company;
|
import ru.clearing.classes.statics.data.company.Company;
|
||||||
import ru.clearing.classes.statics.data.company.relation.Relation;
|
import ru.clearing.classes.statics.data.company.relation.Relation;
|
||||||
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
|
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
|
||||||
|
import ru.clearing.classes.statics.data.misc.STrades;
|
||||||
|
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
|
||||||
|
import ru.clearing.classes.statics.data.security.Security;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
import ru.spcex.clearing.service.validation.ClrngValidationStored;
|
import ru.spcex.clearing.service.validation.ClrngValidationStored;
|
||||||
import ru.spcex.clearing.service.validation.ExecutionDepositValidationRule;
|
import ru.spcex.clearing.service.validation.ExecutionDepositValidationRule;
|
||||||
|
import ru.spcex.clearing.service.validation.STradesValidationRule;
|
||||||
import ru.spcex.platform.classes.base.SpcexObjectBase;
|
import ru.spcex.platform.classes.base.SpcexObjectBase;
|
||||||
import ru.spcex.platform.enumeration.ClearingCategory;
|
import ru.spcex.platform.enumeration.ClearingCategory;
|
||||||
import ru.spcex.platform.imdg.api.Imdg;
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
|
|
@ -22,6 +26,7 @@ import java.util.HashMap;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.function.BiConsumer;
|
import java.util.function.BiConsumer;
|
||||||
import java.util.function.BiFunction;
|
import java.util.function.BiFunction;
|
||||||
|
import java.util.function.Function;
|
||||||
|
|
||||||
@Configuration
|
@Configuration
|
||||||
public class ValidationConfig {
|
public class ValidationConfig {
|
||||||
|
|
@ -36,6 +41,9 @@ public class ValidationConfig {
|
||||||
addImdg.accept(IMDGDistributedNames.Map_Company, Company.class);
|
addImdg.accept(IMDGDistributedNames.Map_Company, Company.class);
|
||||||
addImdg.accept(IMDGDistributedNames.Map_Account, Account.class);
|
addImdg.accept(IMDGDistributedNames.Map_Account, Account.class);
|
||||||
addImdg.accept(IMDGDistributedNames.Map_AccountBalance, AccountBalance.class);
|
addImdg.accept(IMDGDistributedNames.Map_AccountBalance, AccountBalance.class);
|
||||||
|
addImdg.accept(IMDGDistributedNames.Map_Security, Security.class);
|
||||||
|
addImdg.accept(IMDGDistributedNames.Map_Company, Company.class);
|
||||||
|
addImdg.accept(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|
@ -63,4 +71,18 @@ public class ValidationConfig {
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Bean("sTradesValidator")
|
||||||
|
public Function<STrades, IValidator> sTradesValidator() {
|
||||||
|
return sTrades -> {
|
||||||
|
ImdgValidationContext<STrades> context = new ImdgValidationContext<>();
|
||||||
|
context.setValidatedObject(sTrades);
|
||||||
|
addImdg.accept(context, IMDGDistributedNames.Map_Security);
|
||||||
|
addImdg.accept(context, IMDGDistributedNames.Map_Company);
|
||||||
|
addImdg.accept(context, IMDGDistributedNames.Map_TradingClearingRegistry);
|
||||||
|
return new ValidatorImpl<>(context,
|
||||||
|
STradesValidationRule.SecurityPresent,
|
||||||
|
STradesValidationRule.CompanyPresent,
|
||||||
|
STradesValidationRule.TradingClearingRegistryPresent);
|
||||||
|
};
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,13 +1,17 @@
|
||||||
package ru.spcex.clearing.error;
|
package ru.spcex.clearing.error;
|
||||||
|
|
||||||
import ru.spcex.platform.utils.enumeration.IEnumId;
|
import ru.spcex.platform.utils.enumeration.IErrorEnumId;
|
||||||
|
|
||||||
public enum ClearingError implements IEnumId {
|
public enum ClearingError implements IErrorEnumId {
|
||||||
GeneralError(5400L),
|
GeneralError(5400L),
|
||||||
RecordNotFound(5406L),
|
RecordNotFound(5406L),
|
||||||
CompanyCreditCheck(5412L),
|
CompanyCreditCheck(5412L),
|
||||||
CompanyDebitCheck(5413L),
|
CompanyDebitCheck(5413L),
|
||||||
CompanyNotFound(5410L),
|
CompanyNotFound(5410L),
|
||||||
|
SecurityNotFound(5416L),
|
||||||
|
TradingClearingRegistryNotFound(5418L),
|
||||||
|
TradingClearingRegistryNotActive(5419L),
|
||||||
|
NewDealsNotFound(5423L),
|
||||||
;
|
;
|
||||||
private final Long id;
|
private final Long id;
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1,11 +1,10 @@
|
||||||
package ru.spcex.clearing.service;
|
package ru.spcex.clearing.service;
|
||||||
|
|
||||||
import org.apache.kafka.clients.producer.Producer;
|
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.Logger;
|
||||||
import org.slf4j.LoggerFactory;
|
import org.slf4j.LoggerFactory;
|
||||||
import org.springframework.beans.factory.annotation.Autowired;
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
|
import org.springframework.beans.factory.annotation.Qualifier;
|
||||||
import org.springframework.scheduling.annotation.EnableScheduling;
|
import org.springframework.scheduling.annotation.EnableScheduling;
|
||||||
import org.springframework.scheduling.annotation.Scheduled;
|
import org.springframework.scheduling.annotation.Scheduled;
|
||||||
import org.springframework.stereotype.Component;
|
import org.springframework.stereotype.Component;
|
||||||
|
|
@ -14,33 +13,31 @@ import ru.clearing.classes.statics.data.company.Company;
|
||||||
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
|
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
|
||||||
import ru.clearing.classes.statics.data.misc.Listing;
|
import ru.clearing.classes.statics.data.misc.Listing;
|
||||||
import ru.clearing.classes.statics.data.misc.STrades;
|
import ru.clearing.classes.statics.data.misc.STrades;
|
||||||
|
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
|
||||||
import ru.clearing.classes.statics.data.security.Security;
|
import ru.clearing.classes.statics.data.security.Security;
|
||||||
import ru.spcex.clearing.error.ClearingError;
|
import ru.spcex.clearing.error.ClearingError;
|
||||||
import ru.spcex.clearing.error.ClearingException;
|
import ru.spcex.clearing.error.ClearingException;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
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.Consts;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.registry.DealRegisterNewRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.registry.DealRegisterNewRequest;
|
||||||
import ru.spcex.platform.enumeration.Allowed;
|
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||||
import ru.spcex.platform.enumeration.Market;
|
import ru.spcex.clearing.service.validation.ValidationStored;
|
||||||
import ru.spcex.platform.enumeration.MoneyFlowSide;
|
import ru.spcex.platform.enumeration.*;
|
||||||
import ru.spcex.platform.imdg.api.Imdg;
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
import ru.spcex.platform.imdg.api.ImdgId;
|
import ru.spcex.platform.imdg.api.ImdgId;
|
||||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
|
import ru.spcex.platform.utils.enumeration.*;
|
||||||
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
|
|
||||||
import ru.spcex.platform.utils.enumeration.EnumMessage;
|
|
||||||
import ru.spcex.platform.utils.log.ExceptionUtils;
|
import ru.spcex.platform.utils.log.ExceptionUtils;
|
||||||
import ru.spcex.platform.utils.time.TimeUtil;
|
import ru.spcex.platform.utils.time.TimeUtil;
|
||||||
|
import ru.spcex.platform.utils.validation.IValidator;
|
||||||
|
|
||||||
import java.math.BigDecimal;
|
import java.math.BigDecimal;
|
||||||
import java.time.Instant;
|
import java.time.Instant;
|
||||||
import java.time.LocalDate;
|
import java.time.LocalDate;
|
||||||
import java.util.*;
|
import java.util.Collection;
|
||||||
import java.util.concurrent.ExecutionException;
|
import java.util.Map;
|
||||||
import java.util.concurrent.Future;
|
import java.util.Optional;
|
||||||
import java.util.stream.Collectors;
|
import java.util.function.Function;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 1.35. executionDeposit - Сделки
|
* 1.35. executionDeposit - Сделки
|
||||||
|
|
@ -64,10 +61,15 @@ public class ExecutionDepositComponent {
|
||||||
|
|
||||||
Long tradeNum;
|
Long tradeNum;
|
||||||
Instant tradingDay;
|
Instant tradingDay;
|
||||||
|
private final IMessageResolver msgResolver = new SimpleMessageResolver();
|
||||||
|
private final Function<STrades, IValidator> stradesValidator;
|
||||||
|
private final KafkaSender kafkaSender;
|
||||||
|
|
||||||
|
|
||||||
@Autowired
|
@Autowired
|
||||||
public ExecutionDepositComponent(ImdgProvider imdgProvider, Producer<String, Object> kafka) {
|
public ExecutionDepositComponent(ImdgProvider imdgProvider, Producer<String, Object> kafka,
|
||||||
|
@Qualifier("sTradesValidator") Function<STrades, IValidator> stradesValidator,
|
||||||
|
@Qualifier("kafkaSenderWithoutRequestInfo") KafkaSender kafkaSender) {
|
||||||
this.imdgProvider = imdgProvider;
|
this.imdgProvider = imdgProvider;
|
||||||
this.sTradeImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_STrades, STrades.class);
|
this.sTradeImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_STrades, STrades.class);
|
||||||
this.securityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Security, Security.class);
|
this.securityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Security, Security.class);
|
||||||
|
|
@ -78,6 +80,8 @@ public class ExecutionDepositComponent {
|
||||||
|
|
||||||
this.idGenerator = imdgProvider.getImdgIdGenerator();
|
this.idGenerator = imdgProvider.getImdgIdGenerator();
|
||||||
this.kafka = kafka;
|
this.kafka = kafka;
|
||||||
|
this.stradesValidator = stradesValidator;
|
||||||
|
this.kafkaSender = kafkaSender;
|
||||||
|
|
||||||
resetTradingDay();
|
resetTradingDay();
|
||||||
}
|
}
|
||||||
|
|
@ -98,102 +102,53 @@ public class ExecutionDepositComponent {
|
||||||
|
|
||||||
public void processNewTS() {
|
public void processNewTS() {
|
||||||
log.debug("Start check new S_TRADE after {}", tradingDay);
|
log.debug("Start check new S_TRADE after {}", tradingDay);
|
||||||
Collection<STrades> sTrades;
|
//выбираем STrades на сегодня с правильным section
|
||||||
{
|
Collection<STrades> sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of(
|
||||||
ImdgPredicateBuilder pb = sTradeImdg.predicateBuilder();
|
"tradeDate", LocalDate.now(),
|
||||||
ImdgPredicate sql = pb.greatEqual("tradeDateTime", tradingDay);
|
"section", Section.MKR.getKey()
|
||||||
sTrades = sTradeImdg.getCollectionObjectsByPredicate(sql);
|
));
|
||||||
}
|
log.info("Found {} s_trade for today", sTrades.size());
|
||||||
log.info("Found {} new s_trade with trade_num>{}", sTrades.size(), tradeNum);
|
|
||||||
|
|
||||||
if (sTrades.isEmpty()) {
|
if (sTrades.isEmpty()) {
|
||||||
log.info("No new sTrades.");
|
logError(ClearingError.NewDealsNotFound);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
// Выявление новых сделок необходимо выполнить следующие контрольные проверки:
|
//убираем уже добавленные в ExecutionDeposit
|
||||||
|
sTrades.removeIf(sTrd -> {
|
||||||
// Проверить все инструменты.
|
Side sTrdSide = IEnumKey.getEnumByKey(Side.class, sTrd.getOperation());
|
||||||
Set<String> secCodesOfSecurity;
|
if (sTrdSide == null || !(sTrdSide.equals(Side.BUY) || sTrdSide.equals(Side.SELL))) {
|
||||||
{
|
log.warn("cannot parse strade.tradeNum={} strade.operation {}", sTrd.getTradeNum(), sTrd.getOperation());
|
||||||
Set<String> secCodesOfSTrade = sTrades.stream().map(STrades::getSecCode).filter(Objects::nonNull).collect(Collectors.toSet());
|
return true;
|
||||||
log.debug("Verify {} instruments for {} STrade's.", secCodesOfSTrade.size(), sTrades.size());
|
|
||||||
ImdgPredicate allIn = securityImdg.predicateBuilder().in("securitySymbol", secCodesOfSTrade.toArray(new String[0]));
|
|
||||||
Collection<Security> foundSecurities = securityImdg.getCollectionObjectsByPredicate(allIn);
|
|
||||||
secCodesOfSecurity = foundSecurities.stream().map(Security::getSecuritySymbol).filter(Objects::nonNull).collect(Collectors.toSet());
|
|
||||||
if (secCodesOfSecurity.containsAll(secCodesOfSTrade)) {
|
|
||||||
log.debug("All {} Security found by {} secCodes from STrade",
|
|
||||||
secCodesOfSecurity.size(), secCodesOfSTrade.size());
|
|
||||||
} else {
|
|
||||||
HashSet<String> notFoundSymbol = new HashSet<>(secCodesOfSTrade);
|
|
||||||
notFoundSymbol.removeAll(secCodesOfSecurity);
|
|
||||||
log.info("Found only {} Security by {} secCodes from STrade. Not found: {}",
|
|
||||||
secCodesOfSecurity.size(), secCodesOfSTrade.size(), notFoundSymbol);
|
|
||||||
}
|
}
|
||||||
}
|
MoneyFlowSide excDepSide = sTrdSide.equals(Side.BUY) ? MoneyFlowSide.BUY : MoneyFlowSide.SELL;
|
||||||
|
return executionDepositImdg.getSingleObjectByFieldValues(
|
||||||
|
Map.of("tradingDate", sTrd.getTradeDate(),
|
||||||
|
"exchangeExecutionId", sTrd.getTradeNum(),
|
||||||
|
"side", excDepSide.getKey())) != null;
|
||||||
|
});
|
||||||
|
log.info("{} strades left after already-added filtering", sTrades.size());
|
||||||
|
|
||||||
// Необходимо проверять отсутствие ExecutionDeposit с exchangeExecutionId и clearingDate и side.
|
for (STrades sTrd : sTrades) {
|
||||||
// Но при этом, во время создания новых ExecutionDeposit по STrade, TradeNum могут повторяться.
|
log.trace("S_TRADE[{}] new", sTrd.getId());
|
||||||
LocalDate today = LocalDate.now();
|
|
||||||
|
|
||||||
for (STrades trade : sTrades) {
|
//проверки
|
||||||
log.trace("Check s_trade[{}].tradeNum={} operation={} on date {}", trade.getId(), trade.getTradeNum(), trade.getOperation(), today);
|
IValidator validator = stradesValidator.apply(sTrd);
|
||||||
boolean existEDeposit;
|
Optional<EnumMessage> error = validator.tillFirstError();
|
||||||
{
|
if (error.isPresent()) {
|
||||||
ImdgPredicateBuilder pb = executionDepositImdg.predicateBuilder();
|
logError(sTrd.getId(), error.get());
|
||||||
|
continue;
|
||||||
ImdgPredicate predicateOperation = pb.equals("side", trade.getOperation());
|
|
||||||
if ("B".equalsIgnoreCase(trade.getOperation())) {
|
|
||||||
predicateOperation = pb.or(predicateOperation, pb.equals("side", MoneyFlowSide.BUY.getKey()));
|
|
||||||
}
|
|
||||||
// if ("BUY".equalsIgnoreCase(trade.getOperation())) {
|
|
||||||
// predicateOperation = pb.or(predicateOperation, pb.equals("side", "B"));
|
|
||||||
// }
|
|
||||||
if ("S".equalsIgnoreCase(trade.getOperation())) {
|
|
||||||
predicateOperation = pb.or(predicateOperation, pb.equals("side", MoneyFlowSide.SELL.getKey()));
|
|
||||||
}
|
|
||||||
// if ("SELL".equalsIgnoreCase(trade.getOperation())) {
|
|
||||||
// predicateOperation = pb.or(predicateOperation, pb.equals("side", "S"));
|
|
||||||
// }
|
|
||||||
|
|
||||||
ImdgPredicate predicate = pb.and(
|
|
||||||
pb.equals("exchangeExecutionId", trade.getTradeNum()),
|
|
||||||
pb.equals("clearingDate", today),
|
|
||||||
predicateOperation
|
|
||||||
);
|
|
||||||
Collection<ExecutionDeposit> existsEDeposit = executionDepositImdg.getCollectionObjectsByPredicate(predicate);
|
|
||||||
// Collection<ExecutionDeposit> existsEDeposit=executionDepositImdg.getCollectionObjectsByFieldValues(Map.of(
|
|
||||||
// "exchangeExecutionId", trade.getTradeNum(),
|
|
||||||
// "clearingDate", today
|
|
||||||
// ));
|
|
||||||
existEDeposit = !existsEDeposit.isEmpty();
|
|
||||||
log.trace("Check exist ExecutionDeposit, SQL \"{}\", found {}, exist {}", predicate, existsEDeposit.size(), existEDeposit);
|
|
||||||
}
|
}
|
||||||
|
ExecutionDeposit newED;
|
||||||
if (!existEDeposit) {
|
try {
|
||||||
log.trace("S_TRADE[{}] new", trade.getId());
|
newED = createExecutionDeposit(sTrd, validator);
|
||||||
|
executionDepositImdg.insert(newED);
|
||||||
// Проверка secCode
|
sendNotification(newED);
|
||||||
if (!secCodesOfSecurity.contains(trade.getSecCode())) {
|
log.debug("New executionDeposit.id={} was created.", newED.getId());
|
||||||
log.error("Error {}: STrade[{}].secCode={} not found",
|
} catch (ClearingException ce) {
|
||||||
ClearingError.RecordNotFound.getId(), trade.getId(), trade.getSecCode());
|
auditMessage(ce);
|
||||||
continue;
|
} catch (Exception e) {
|
||||||
}
|
log.error("When create new ExecutionDeposit by STrade[{}] error: {}", sTrd.getId(), ExceptionUtils.getStackTrace(e));
|
||||||
|
|
||||||
ExecutionDeposit newED;
|
|
||||||
try {
|
|
||||||
newED = createExecutionDeposit(trade, null, null);
|
|
||||||
executionDepositImdg.insert(newED);
|
|
||||||
sendNotification(newED);
|
|
||||||
log.debug("New executionDeposit.id={} was created.", newED.getId());
|
|
||||||
} catch (ClearingException ce) {
|
|
||||||
auditMessage(ce);
|
|
||||||
} catch (Exception e) {
|
|
||||||
log.error("When create new ExecutionDeposit by STrade[{}] error: {}", trade.getId(), ExceptionUtils.getStackTrace(e));
|
|
||||||
}
|
|
||||||
|
|
||||||
} else {
|
|
||||||
log.trace("S_TRADE[{}] already has executionDeposit: exchangeExecutionId={} on date {}", trade.getId(), trade.getTradeNum(), today);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -207,106 +162,79 @@ public class ExecutionDepositComponent {
|
||||||
log.error("AUDIT error code {}: {}", ce.getEnumMsg(), ce.getMessage());
|
log.error("AUDIT error code {}: {}", ce.getEnumMsg(), ce.getMessage());
|
||||||
}
|
}
|
||||||
|
|
||||||
protected void auditMessage(String message, Object... ids) {
|
protected void sendNotification(ExecutionDeposit forED) throws ClearingException {
|
||||||
String txt = message;
|
|
||||||
if (ids != null && ids.length > 0) {
|
|
||||||
txt += " object:" + Arrays.toString(ids);
|
|
||||||
}
|
|
||||||
log.error("audit \"clearing-service\", errorText: {}", txt);
|
|
||||||
}
|
|
||||||
|
|
||||||
protected void sendNotification(ExecutionDeposit forED) {
|
|
||||||
final String destination = Consts.REGISTRY_DEAL_REGISTER_NEW;
|
final String destination = Consts.REGISTRY_DEAL_REGISTER_NEW;
|
||||||
DealRegisterNewRequest requestPayload = new DealRegisterNewRequest();
|
DealRegisterNewRequest requestPayload = new DealRegisterNewRequest();
|
||||||
requestPayload.setExecutionId(forED.getId());
|
requestPayload.setExecutionId(forED.getId());
|
||||||
requestPayload.setExchangeExecutionId(forED.getExchangeExecutionId());
|
requestPayload.setExchangeExecutionId(forED.getExchangeExecutionId());
|
||||||
// requestPayload.setId(idGenerator.nextId());
|
Long generatedRequestId = kafkaSender.sendRequestToQueue(destination, requestPayload);
|
||||||
|
if (generatedRequestId == null) {
|
||||||
BaseRequest<Object> request = new BaseRequest<>();
|
throw new RuntimeException("failed to send request to " + destination);
|
||||||
request.setId(idGenerator.nextId());
|
|
||||||
request.setActionType(ActionType.NEW);
|
|
||||||
request.setRequestPayload(requestPayload);
|
|
||||||
log.trace("Send to {} new ExecutionDeposit[{}]", destination, forED.getId());
|
|
||||||
Future<RecordMetadata> send = kafka.send(new ProducerRecord<>(destination, request));
|
|
||||||
try {
|
|
||||||
send.get();
|
|
||||||
} catch (InterruptedException e) {
|
|
||||||
Thread.currentThread().interrupt();
|
|
||||||
throw new RuntimeException(e);
|
|
||||||
} catch (ExecutionException e) {
|
|
||||||
throw new RuntimeException(e);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
protected ExecutionDeposit createExecutionDeposit(STrades sTrades,
|
protected ExecutionDeposit createExecutionDeposit(STrades sTrades,
|
||||||
Allowed coverageStatus, Long sessionId) throws ClearingException {
|
IValidator validator) throws ClearingException {
|
||||||
Account account = null; // todo CLS-275 accountImdg.getSingleObjectByFieldValues(Map.of("account", sTrades.getMoneyAccount()));
|
Security security = validator.getStored(ValidationStored.STradesSecurity);
|
||||||
Security security = securityImdg.getSingleObjectByFieldValues(Map.of("securitySymbol", sTrades.getSecCode()));
|
Company company = validator.getStored(ValidationStored.STradesCompany);
|
||||||
if (security == null) {
|
Company counterCompany = validator.getStored(ValidationStored.STradesCounterCompany);
|
||||||
log.warn("security securitySymbol=\"{}\" not found", sTrades.getSecCode());
|
TradingClearingRegistry rgstr = validator.getStored(ValidationStored.STradesTradingClearingRegistry);
|
||||||
throw new ClearingException(new EnumMessage(ClearingError.RecordNotFound, sTrades.getSecCode()));
|
Listing listing = listingImdg.getSingleObjectByFieldValues(Map.of("securityId", security.getId()));
|
||||||
}
|
|
||||||
Listing listing = null;
|
|
||||||
if (security != null) {
|
|
||||||
listing = listingImdg.getSingleObjectByFieldValues(Map.of("securityId", security.getId()));
|
|
||||||
}
|
|
||||||
Company company = companyImdg.getSingleObjectByFieldValues(Map.of("tradingCode", sTrades.getFirmId()));
|
|
||||||
if (company == null) {
|
|
||||||
throw new ClearingException(new EnumMessage(ClearingError.CompanyNotFound, sTrades.getFirmId()));
|
|
||||||
}
|
|
||||||
|
|
||||||
return createExecutionDeposit(sTrades, account, listing, company, security, coverageStatus, sessionId);
|
|
||||||
}
|
|
||||||
|
|
||||||
private ExecutionDeposit createExecutionDeposit(STrades sTrades, Account account, Listing listing,
|
|
||||||
Company company, Security security,
|
|
||||||
Allowed coverageStatus, Long sessionId) {
|
|
||||||
ExecutionDeposit eDeposit = new ExecutionDeposit();
|
ExecutionDeposit eDeposit = new ExecutionDeposit();
|
||||||
eDeposit.setId(idGenerator.nextId());
|
|
||||||
final Instant now = Instant.now();
|
final Instant now = Instant.now();
|
||||||
final LocalDate nowDay = TimeUtil.toLocalDate(now);
|
|
||||||
eDeposit.setCreated(now);
|
eDeposit.setCreated(now);
|
||||||
eDeposit.setTradingDate(nowDay);
|
eDeposit.setTradingDate(sTrades.getTradeDate());
|
||||||
eDeposit.setClearingDate(nowDay);
|
eDeposit.setClearingDate(TimeUtil.toLocalDate(now));
|
||||||
|
|
||||||
eDeposit.setExchangeExecutionId(sTrades.getTradeNum());
|
eDeposit.setExchangeExecutionId(sTrades.getTradeNum());
|
||||||
eDeposit.setExchangeExecutionTime(sTrades.getTradeDateTime());
|
eDeposit.setExchangeExecutionTime(sTrades.getTradeDateTime());
|
||||||
if (account != null) {
|
eDeposit.setTradingClearingRegistryId(rgstr.getId());
|
||||||
// todo CLS-275 eDeposit.setAccountId(account.getId());
|
|
||||||
}
|
|
||||||
eDeposit.setMarket(Market.mkrs.getKey());
|
eDeposit.setMarket(Market.mkrs.getKey());
|
||||||
eDeposit.setPrice(sTrades.getPrice());
|
eDeposit.setPrice(sTrades.getPrice());
|
||||||
eDeposit.setLots(sTrades.getQty());
|
eDeposit.setLots(sTrades.getQty());
|
||||||
if (listing != null && listing.getLotSize() != null && eDeposit.getLots() != null) {
|
if (listing != null && listing.getLotSize() != null && eDeposit.getLots() != null) {
|
||||||
BigDecimal quantity = eDeposit.getLots().multiply(listing.getLotSize());
|
BigDecimal lots = eDeposit.getLots();
|
||||||
eDeposit.setQuantity(quantity);
|
BigDecimal listingLotSize = listing.getLotSize();
|
||||||
|
if (listingLotSize != null && lots != null) {
|
||||||
|
eDeposit.setQuantity(lots.multiply(listingLotSize));
|
||||||
|
}
|
||||||
}
|
}
|
||||||
eDeposit.setFirstLegAmount(sTrades.getValue());
|
eDeposit.setFirstLegAmount(sTrades.getValue());
|
||||||
eDeposit.setSecondLegAmount(sTrades.getValue());
|
eDeposit.setSecondLegAmount(sTrades.getValue());
|
||||||
//eDeposit.setInterestAmount(null);
|
{
|
||||||
String operation = sTrades.getOperation(); //Символьный код по справочнику moneyFlowSide), соответствующий значению из s_trade.operation (sTrade.getOperation())
|
Side sTradeSide = IEnumKey.getEnumByKey(Side.class, sTrades.getOperation());
|
||||||
if ("B".equalsIgnoreCase(operation)) {
|
if (sTradeSide != null) {
|
||||||
operation = MoneyFlowSide.BUY.getKey();
|
switch (sTradeSide) {
|
||||||
|
case BUY -> eDeposit.setSide(MoneyFlowSide.BUY.getKey());
|
||||||
|
case SELL -> eDeposit.setSide(MoneyFlowSide.SELL.getKey());
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
if ("S".equalsIgnoreCase(operation)) {
|
eDeposit.setSettlementCurrency(CurrencyCode.RUB.getKey());
|
||||||
operation = MoneyFlowSide.SELL.getKey();
|
|
||||||
}
|
|
||||||
eDeposit.setSide(operation);
|
|
||||||
eDeposit.setSettlementCurrency("RUB"); // (справочник currencyCode)
|
|
||||||
eDeposit.setCompanyId(company.getId());
|
eDeposit.setCompanyId(company.getId());
|
||||||
// todo CLS-275 eDeposit.setDuration(sTrades.getDaysToMatDate());
|
eDeposit.setDuration(sTrades.getRepoTerm());
|
||||||
eDeposit.setFirstLegSettlementDate(nowDay);
|
eDeposit.setFirstLegSettlementDate(sTrades.getTradeDate());
|
||||||
eDeposit.setSecondLegSettlementDate(sTrades.getSettleDate());
|
{
|
||||||
//eDeposit.setFirstLegSettlementCode(null);
|
LocalDate sTrdSettleDate = sTrades.getSettleDate();
|
||||||
//eDeposit.setSecondLegSettlementCode(null);
|
if (sTrdSettleDate != null) {
|
||||||
|
eDeposit.setSecondLegSettlementDate(sTrdSettleDate.plusDays(sTrades.getRepoTerm()));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
eDeposit.setFirstLegSettlementCode(sTrades.getSettleCode());
|
||||||
|
eDeposit.setSecondLegSettlementCode(sTrades.getSettleCode());
|
||||||
eDeposit.setSecurityFullName(security.getFullName());
|
eDeposit.setSecurityFullName(security.getFullName());
|
||||||
eDeposit.setSecuritySymbol(security.getSecuritySymbol());
|
eDeposit.setSecuritySymbol(security.getSecuritySymbol());
|
||||||
eDeposit.setSecurityId(security.getId());
|
eDeposit.setSecurityId(security.getId());
|
||||||
//eDeposit.setCounterPartyId(null);
|
eDeposit.setContract(sTrades.getClassCode());
|
||||||
eDeposit.setCoverageStatus(coverageStatus == null ? null : coverageStatus.getKey()); // Заполняется по справочнику allowed в результате расчета требований и обязательств. TODO
|
eDeposit.setCounterPartyId(counterCompany.getId());
|
||||||
eDeposit.setSessionId(sessionId);
|
|
||||||
|
|
||||||
return eDeposit;
|
return eDeposit;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private void logError(IEnumId subject, Object... args) {
|
||||||
|
log.warn("{}", msgResolver.resolve(new EnumMessage(subject, args)));
|
||||||
|
}
|
||||||
|
|
||||||
|
private void logError(Long sTradeId, EnumMessage msg) {
|
||||||
|
log.warn("sTrade id={} {}", sTradeId, msgResolver.resolve(msg));
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,111 @@
|
||||||
|
package ru.spcex.clearing.service.validation;
|
||||||
|
|
||||||
|
import ru.clearing.classes.statics.data.company.Company;
|
||||||
|
import ru.clearing.classes.statics.data.misc.STrades;
|
||||||
|
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
|
||||||
|
import ru.clearing.classes.statics.data.security.Security;
|
||||||
|
import ru.spcex.clearing.error.ClearingError;
|
||||||
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
import ru.spcex.platform.enumeration.ServiceStatus;
|
||||||
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
|
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
|
||||||
|
import ru.spcex.platform.utils.enumeration.EnumMessage;
|
||||||
|
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
||||||
|
import ru.spcex.platform.utils.text.TextUtil;
|
||||||
|
import ru.spcex.platform.utils.validation.IValidationRule;
|
||||||
|
|
||||||
|
import java.util.Map;
|
||||||
|
import java.util.Optional;
|
||||||
|
|
||||||
|
public enum STradesValidationRule implements IValidationRule<ImdgValidationContext<STrades>> {
|
||||||
|
|
||||||
|
// 2.1. Найти в security запись, у которой security.securitySymbol=sTrades.secCode.
|
||||||
|
// Если такой записи нет, записать в лог ошибку (5416) "Инструмент %s не найден".
|
||||||
|
SecurityPresent() {
|
||||||
|
@Override
|
||||||
|
public Optional<EnumMessage> validate(ImdgValidationContext<STrades> context) {
|
||||||
|
STrades validatedObject = context.getValidatedObject();
|
||||||
|
if (TextUtil.isEmpty(validatedObject.getSecCode())) {
|
||||||
|
return of(ClearingError.SecurityNotFound, validatedObject.getSecCode());
|
||||||
|
}
|
||||||
|
Imdg<Security> securityImdg = context.obtainMap(IMDGDistributedNames.Map_Security, Security.class);
|
||||||
|
Security security = securityImdg.getSingleObjectByFieldValues(Map.of("securitySymbol", validatedObject.getSecCode()));
|
||||||
|
if (security == null) {
|
||||||
|
return of(ClearingError.SecurityNotFound, validatedObject.getSecCode());
|
||||||
|
}
|
||||||
|
context.storeObject(ValidationStored.STradesSecurity, security);
|
||||||
|
return empty();
|
||||||
|
}
|
||||||
|
},
|
||||||
|
// 2.2. Найти в company запись, у которой company.tradingCode=sTrades.firmId.
|
||||||
|
// Если такой записи нет, записать в лог ошибку (5410) "Компания %s не найдена".
|
||||||
|
CompanyPresent() {
|
||||||
|
@Override
|
||||||
|
public Optional<EnumMessage> validate(ImdgValidationContext<STrades> context) {
|
||||||
|
STrades validatedObject = context.getValidatedObject();
|
||||||
|
if (TextUtil.isEmpty(validatedObject.getFirmId())) {
|
||||||
|
return of(ClearingError.CompanyNotFound, validatedObject.getFirmId());
|
||||||
|
}
|
||||||
|
Imdg<Company> companyImdg = context.obtainMap(IMDGDistributedNames.Map_Company, Company.class);
|
||||||
|
Company company = companyImdg.getSingleObjectByFieldValues(Map.of("tradingCode", validatedObject.getFirmId()));
|
||||||
|
if (company == null) {
|
||||||
|
return of(ClearingError.CompanyNotFound, validatedObject.getFirmId());
|
||||||
|
}
|
||||||
|
context.storeObject(ValidationStored.STradesCompany, company);
|
||||||
|
return empty();
|
||||||
|
}
|
||||||
|
},
|
||||||
|
CounterCompanyPresent() {
|
||||||
|
@Override
|
||||||
|
public Optional<EnumMessage> validate(ImdgValidationContext<STrades> context) {
|
||||||
|
STrades validatedObject = context.getValidatedObject();
|
||||||
|
if (TextUtil.isEmpty(validatedObject.getCpFirmId())) {
|
||||||
|
return of(ClearingError.CompanyNotFound, validatedObject.getCpFirmId());
|
||||||
|
}
|
||||||
|
Imdg<Company> companyImdg = context.obtainMap(IMDGDistributedNames.Map_Company, Company.class);
|
||||||
|
Company company = companyImdg.getSingleObjectByFieldValues(Map.of("tradingCode", validatedObject.getCpFirmId()));
|
||||||
|
if (company == null) {
|
||||||
|
return of(ClearingError.CompanyNotFound, validatedObject.getCpFirmId());
|
||||||
|
}
|
||||||
|
context.storeObject(ValidationStored.STradesCounterCompany, company);
|
||||||
|
return empty();
|
||||||
|
}
|
||||||
|
},
|
||||||
|
// 2.3. Найти в tradingClearingRegistry запись, у которой tradingClearingRegistry.code=sTrades.account
|
||||||
|
// и tradingClearingRegistry.companyId=companyId, найденому в предыдущем пункте.
|
||||||
|
// Если такой записи нет, записать в лог ошибку (5418) "Торгово-клиринговый регистр %s не найден для компании %s".
|
||||||
|
TradingClearingRegistryPresent() {
|
||||||
|
@Override
|
||||||
|
public Optional<EnumMessage> validate(ImdgValidationContext<STrades> context) {
|
||||||
|
STrades validatedObject = context.getValidatedObject();
|
||||||
|
if (TextUtil.isEmpty(validatedObject.getAccount())) {
|
||||||
|
return of(ClearingError.TradingClearingRegistryNotFound, validatedObject.getAccount());
|
||||||
|
}
|
||||||
|
if (context.getStoredObject(ValidationStored.STradesCompany) == null) { //maybe unnecessary
|
||||||
|
return of(ClearingError.CompanyNotFound, validatedObject.getSecCode());
|
||||||
|
}
|
||||||
|
Imdg<TradingClearingRegistry> registryImdg = context.obtainMap(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
|
||||||
|
TradingClearingRegistry tcRegister = registryImdg.getSingleObjectByFieldValues(
|
||||||
|
Map.of("code", validatedObject.getAccount(),
|
||||||
|
"companyId", ((Company) context.getStoredObject(ValidationStored.STradesCompany)).getId()
|
||||||
|
));
|
||||||
|
if (tcRegister == null) {
|
||||||
|
return of(ClearingError.TradingClearingRegistryNotFound, validatedObject.getAccount());
|
||||||
|
}
|
||||||
|
// 2.4. Проверить, что найденный торгово-клиринговый регистр в активном состоянии:
|
||||||
|
// tradingClearingRegistry.status≠BLKD/SSPD/CLOS (см. справочник serviceStatus).
|
||||||
|
// Иначе записать в лог ошибку (5419) "Торгово-клиринговый регистр %s неактивен".
|
||||||
|
ServiceStatus status = IEnumKey.getEnumByKey(ServiceStatus.class, tcRegister.getStatus());
|
||||||
|
if (IEnumKey.contains(status, ServiceStatus.Blocked, ServiceStatus.Suspended, ServiceStatus.Closed)) {
|
||||||
|
return of(ClearingError.TradingClearingRegistryNotActive, validatedObject.getAccount());
|
||||||
|
}
|
||||||
|
context.storeObject(ValidationStored.STradesTradingClearingRegistry, tcRegister);
|
||||||
|
return empty();
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public String ruleName() {
|
||||||
|
return "STradesValidationRule." + name();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,5 @@
|
||||||
|
package ru.spcex.clearing.service.validation;
|
||||||
|
|
||||||
|
public enum ValidationStored {
|
||||||
|
STradesCompany, STradesCounterCompany, STradesSecurity, STradesTradingClearingRegistry
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,24 @@
|
||||||
|
package ru.spcex.platform.enumeration;
|
||||||
|
|
||||||
|
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
||||||
|
|
||||||
|
public enum Section implements IEnumKey {
|
||||||
|
MKR("MKR"),
|
||||||
|
FOND("FOND");
|
||||||
|
|
||||||
|
Section(String key) {
|
||||||
|
this.key = key;
|
||||||
|
}
|
||||||
|
|
||||||
|
private String key;
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public String getKey() {
|
||||||
|
return this.key;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean equalsByKey(String key) {
|
||||||
|
return IEnumKey.super.equalsByKey(key);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -24,6 +24,7 @@ public class KafkaSender {
|
||||||
private Producer<String, Object> kafka;
|
private Producer<String, Object> kafka;
|
||||||
private Supplier<Long> idGenerator;
|
private Supplier<Long> idGenerator;
|
||||||
private Function<String, KafkaImdgInsert> imdgGenerator;
|
private Function<String, KafkaImdgInsert> imdgGenerator;
|
||||||
|
private boolean saveRequestInfo = true;
|
||||||
|
|
||||||
KafkaSender() {
|
KafkaSender() {
|
||||||
this.allImdgMaps = new ConcurrentHashMap<>();
|
this.allImdgMaps = new ConcurrentHashMap<>();
|
||||||
|
|
@ -55,6 +56,12 @@ public class KafkaSender {
|
||||||
}
|
}
|
||||||
|
|
||||||
private void saveRequestToStorage(String destination, BaseRequest<Object> request) {
|
private void saveRequestToStorage(String destination, BaseRequest<Object> request) {
|
||||||
|
if (!saveRequestInfo) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (imdgGenerator == null) {
|
||||||
|
throw new IllegalStateException("imdgGenerator is null");
|
||||||
|
}
|
||||||
KafkaImdgInsert imdgInsert = getImdg(destination);
|
KafkaImdgInsert imdgInsert = getImdg(destination);
|
||||||
RequestInfo requestInfo = RequestInfo.create(request.getId());
|
RequestInfo requestInfo = RequestInfo.create(request.getId());
|
||||||
imdgInsert.insert(requestInfo);
|
imdgInsert.insert(requestInfo);
|
||||||
|
|
@ -76,4 +83,8 @@ public class KafkaSender {
|
||||||
void setImdgProvider(Function<String, KafkaImdgInsert> imdgGenerator) {
|
void setImdgProvider(Function<String, KafkaImdgInsert> imdgGenerator) {
|
||||||
this.imdgGenerator = imdgGenerator;
|
this.imdgGenerator = imdgGenerator;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
void setSaveRequestInfo(boolean saveRequestInfo) {
|
||||||
|
this.saveRequestInfo = saveRequestInfo;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -12,5 +12,7 @@ public interface KafkaSenderBuilder {
|
||||||
|
|
||||||
KafkaSenderBuilder imdgProvider(Function<String, KafkaImdgInsert> imdgGenerator);
|
KafkaSenderBuilder imdgProvider(Function<String, KafkaImdgInsert> imdgGenerator);
|
||||||
|
|
||||||
|
KafkaSenderBuilder saveRequestInfo(boolean saveRequestInfo);
|
||||||
|
|
||||||
KafkaSender build();
|
KafkaSender build();
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -31,6 +31,12 @@ public class KafkaSenderBuilderImpl implements KafkaSenderBuilder{
|
||||||
return this;
|
return this;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public KafkaSenderBuilder saveRequestInfo(boolean saveRequestInfo) {
|
||||||
|
kafkaSender.setSaveRequestInfo(saveRequestInfo);
|
||||||
|
return this;
|
||||||
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public KafkaSender build() {
|
public KafkaSender build() {
|
||||||
return kafkaSender;
|
return kafkaSender;
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,78 @@
|
||||||
|
package ru.spcex.clearing.platform.messaging.service.sender;
|
||||||
|
|
||||||
|
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.kafka.requestreply.ReplyingKafkaTemplate;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.ActionType;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
|
||||||
|
import ru.spcex.platform.classes.base.SpcexObjectBase;
|
||||||
|
import ru.spcex.platform.utils.log.ExceptionUtils;
|
||||||
|
|
||||||
|
import java.util.Map;
|
||||||
|
import java.util.concurrent.ConcurrentHashMap;
|
||||||
|
import java.util.concurrent.ExecutionException;
|
||||||
|
import java.util.concurrent.Future;
|
||||||
|
import java.util.function.Function;
|
||||||
|
import java.util.function.Supplier;
|
||||||
|
|
||||||
|
public class KafkaSyncRequestReplySender {
|
||||||
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
private final Map<String, KafkaImdgInsert> allImdgMaps;
|
||||||
|
private Producer<String, Object> kafka;
|
||||||
|
private Supplier<Long> idGenerator;
|
||||||
|
private Function<String, KafkaImdgInsert> imdgGenerator;
|
||||||
|
private ReplyingKafkaTemplate<String, Object, Object> kafkaTemplate;
|
||||||
|
|
||||||
|
KafkaSyncRequestReplySender() {
|
||||||
|
this.allImdgMaps = new ConcurrentHashMap<>();
|
||||||
|
}
|
||||||
|
|
||||||
|
public static KafkaSenderBuilderImpl setup() {
|
||||||
|
return new KafkaSenderBuilderImpl();
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
public Long sendRequestToQueue(String destination, Object requestPayload) {
|
||||||
|
BaseRequest<Object> request = new BaseRequest<>();
|
||||||
|
request.setId(idGenerator.get());
|
||||||
|
request.setActionType(ActionType.SYSTEM);
|
||||||
|
request.setRequestPayload(requestPayload);
|
||||||
|
//сохраняет данные о запросе в хранилище
|
||||||
|
saveRequestToStorage(destination, request);
|
||||||
|
Future<RecordMetadata> send = kafka.send(new ProducerRecord<>(destination, request));
|
||||||
|
try {
|
||||||
|
send.get();
|
||||||
|
} catch (InterruptedException | ExecutionException e) {
|
||||||
|
log.error(ExceptionUtils.getStackTrace(e));
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
return request.getId();
|
||||||
|
}
|
||||||
|
|
||||||
|
private void saveRequestToStorage(String destination, BaseRequest<Object> request) {
|
||||||
|
KafkaImdgInsert imdgInsert = getImdg(destination);
|
||||||
|
RequestInfo requestInfo = RequestInfo.create(request.getId());
|
||||||
|
imdgInsert.insert(requestInfo);
|
||||||
|
}
|
||||||
|
|
||||||
|
@SuppressWarnings("unchecked")
|
||||||
|
private <T extends SpcexObjectBase> KafkaImdgInsert getImdg(String mapName) {
|
||||||
|
return allImdgMaps.computeIfAbsent(mapName, (mapName1) -> imdgGenerator.apply(mapName));
|
||||||
|
}
|
||||||
|
|
||||||
|
void setProducer(Producer<String, Object> kafka) {
|
||||||
|
this.kafka = kafka;
|
||||||
|
}
|
||||||
|
|
||||||
|
void setIdGenerator(Supplier<Long> idGenerator) {
|
||||||
|
this.idGenerator = idGenerator;
|
||||||
|
}
|
||||||
|
|
||||||
|
void setImdgProvider(Function<String, KafkaImdgInsert> imdgGenerator) {
|
||||||
|
this.imdgGenerator = imdgGenerator;
|
||||||
|
}
|
||||||
|
}
|
||||||
Loading…
Add table
Reference in a new issue