clearing-service changeStatusExtract => gateway without timeout
optional max-poll-interval-ms setting
This commit is contained in:
parent
ab80d194a5
commit
0fd9208e41
6 changed files with 61 additions and 29 deletions
|
|
@ -2,6 +2,8 @@ package ru.spcex.clearing.config;
|
|||
|
||||
import org.apache.kafka.clients.consumer.Consumer;
|
||||
import org.apache.kafka.clients.producer.Producer;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
|
|
@ -26,19 +28,31 @@ import java.util.function.Supplier;
|
|||
|
||||
@Configuration
|
||||
public class KafkaConfig {
|
||||
@Autowired
|
||||
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
|
||||
@Bean("kafkaConsumer")
|
||||
public Consumer<String, Object> createConsumer(ClearingServiceSettings settings) {
|
||||
return KafkaConsumerFactory.consumer(settings.getKafkaConsumer());
|
||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||
private final KafkaConsumerSettings kafkaSettings;
|
||||
|
||||
public KafkaConfig(ClearingServiceSettings settings) {
|
||||
this.kafkaSettings = settings.getKafkaConsumer();
|
||||
Integer gtwTimeout = settings.getSessionStage().getInspectionGatewayTimeout();
|
||||
if (gtwTimeout != null) {
|
||||
int maxPollIntervalMs = gtwTimeout * 1000 + 1000;
|
||||
log.debug("setting max.poll.interval.ms {}", maxPollIntervalMs);
|
||||
this.kafkaSettings.setManuallyMaxPollIntervalMs(maxPollIntervalMs);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
|
||||
@Bean("kafkaConsumer")
|
||||
public Consumer<String, Object> createConsumer() {
|
||||
return KafkaConsumerFactory.consumer(kafkaSettings);
|
||||
}
|
||||
|
||||
@Autowired
|
||||
@Bean("kafkaConsumerGateway")
|
||||
public Supplier<Consumer<String, Object>> getwaySessionConsumer(ClearingServiceSettings settings) {
|
||||
public Supplier<Consumer<String, Object>> getwaySessionConsumer() {
|
||||
AtomicInteger groupId = new AtomicInteger(1);
|
||||
return () -> {
|
||||
KafkaConsumerSettings consumer = settings.getKafkaConsumer();
|
||||
KafkaConsumerSettings consumer = kafkaSettings;
|
||||
KafkaConsumerSettings gatewayConsumer = new KafkaConsumerSettings();
|
||||
gatewayConsumer.setBootstrapServers(consumer.getBootstrapServers());
|
||||
gatewayConsumer.setSessionTimeoutMs(consumer.getSessionTimeoutMs());
|
||||
|
|
|
|||
|
|
@ -368,6 +368,7 @@ public class ValidationConfig {
|
|||
new PresentById(IMDGDistributedNames.Map_Registry, ClearingError.RecordNotFound, true),
|
||||
StatusExtractValidationRule.RegistryCodeCheck,
|
||||
StatusExtractValidationRule.ContractCheck,
|
||||
StatusExtractValidationRule.RequestStatusValid,
|
||||
StatusExtractValidationRule.RegistryStatusValid
|
||||
);
|
||||
};
|
||||
|
|
|
|||
|
|
@ -48,16 +48,11 @@ import ru.spcex.platform.utils.validation.IValidator;
|
|||
import java.math.BigDecimal;
|
||||
import java.time.Instant;
|
||||
import java.time.LocalDate;
|
||||
import java.util.Collection;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.*;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.Supplier;
|
||||
import java.util.stream.Collectors;
|
||||
import java.util.stream.Stream;
|
||||
|
||||
import static ru.spcex.clearing.session.stage.impl.GatewayRequester.mapError;
|
||||
import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD;
|
||||
|
||||
@Service
|
||||
|
|
@ -290,23 +285,11 @@ public class RegistryService {
|
|||
return reqHelp.error(req.getId(), err.get());
|
||||
}
|
||||
Registry rgs = validator.getStored(Stored.PresentById);
|
||||
Supplier<AssetOperationRequest> gtwBuilder = () -> GatewayRequestCreator.from(rgs, InOutDirection.in);
|
||||
Optional<Boolean> gatewayOk;
|
||||
if (trdTime.isTradingTime()) {
|
||||
gatewayOk = gateway.gatewayRequestAndWait(gtwBuilder);
|
||||
Long gtwReqId = kafkaSender.sendRequestToQueue(Consts.ASSET_OPERATION, gtwReq(rgs));
|
||||
log.debug("send request to {} id={}", Consts.ASSET_OPERATION, gtwReqId);
|
||||
} else {
|
||||
log.debug("RegistryChangeStatusExtractRequest rgs.id={} not sending gateway request", rgs.getId());
|
||||
gatewayOk = Optional.of(true);
|
||||
}
|
||||
if (gatewayOk.isEmpty() || !gatewayOk.get()) {
|
||||
String gtwErr = msgResolver.resolve(mapError(gatewayOk));
|
||||
log.error("{}.id={} {}",
|
||||
rgs.getRegistryCode(),
|
||||
rgs.getId(),
|
||||
gtwErr);
|
||||
notification.sendNotification(ObjectType.rgst,
|
||||
"Отметка о получении выписки: %s".formatted(gtwErr),
|
||||
Priority.HIGH);
|
||||
}
|
||||
rgs.setRegistryStatus(payload.getRegistryStatus());
|
||||
log.debug("changing registry.id={} status to {}", rgs.getId(), payload.getRegistryStatus());
|
||||
|
|
@ -321,6 +304,13 @@ public class RegistryService {
|
|||
return reqHelp.success(req.getId());
|
||||
}
|
||||
|
||||
private AssetOperationListRequest gtwReq(Registry rgs) {
|
||||
AssetOperationRequest item = GatewayRequestCreator.from(rgs, InOutDirection.in);
|
||||
AssetOperationListRequest gtwReq = new AssetOperationListRequest();
|
||||
gtwReq.setAssetOperationRequests(Collections.singletonList(item));
|
||||
return gtwReq;
|
||||
}
|
||||
|
||||
public RequestInfoUpdate identificationFunds(BaseRequest<IdentificationFundsRequest> req) {
|
||||
if (!rights.userHasRole(req.getUserId(), UserRole.Admin)) {
|
||||
return reqHelp.error(req.getId(), ClearingError.UserVerifyDenial);
|
||||
|
|
|
|||
|
|
@ -63,7 +63,7 @@ public enum StatusExtractValidationRule implements IValidationRule<ImdgValidatio
|
|||
return empty();
|
||||
}
|
||||
},
|
||||
RegistryStatusValid() {
|
||||
RequestStatusValid() {
|
||||
@Override
|
||||
public Optional<EnumMessage> validate(ImdgValidationContext<RegistryChangeStatusExtractRequest> context) {
|
||||
RegistryChangeStatusExtractRequest validatedObject = context.getValidatedObject();
|
||||
|
|
@ -76,6 +76,17 @@ public enum StatusExtractValidationRule implements IValidationRule<ImdgValidatio
|
|||
}
|
||||
return empty();
|
||||
}
|
||||
},
|
||||
RegistryStatusValid() {
|
||||
@Override
|
||||
public Optional<EnumMessage> validate(ImdgValidationContext<RegistryChangeStatusExtractRequest> context) {
|
||||
Registry rgs = context.getStoredObject(Stored.PresentById);
|
||||
RegistryStatus status = IEnumKey.getEnumByKey(RegistryStatus.class, rgs.getRegistryStatus());
|
||||
if (RegistryStatus.OK.equals(status)) {
|
||||
return of(ClearingError.ObligationsAlreadyCalculated);
|
||||
}
|
||||
return empty();
|
||||
}
|
||||
}
|
||||
;
|
||||
private final static Logger log = LoggerFactory.getLogger(StatusExtractValidationRule.class);
|
||||
|
|
|
|||
|
|
@ -22,6 +22,9 @@ public class KafkaConsumerFactory {
|
|||
props.put("enable.auto.commit", kafkaSettings.getEnableAutoCommit().toString());
|
||||
props.put("session.timeout.ms", kafkaSettings.getSessionTimeoutMs().toString());
|
||||
props.put("auto.offset.reset", kafkaSettings.getAutoOffsetReset());
|
||||
if (kafkaSettings.getMaxPollIntervalMs() != null) {
|
||||
props.put("max.poll.interval.ms", String.valueOf(kafkaSettings.getMaxPollIntervalMs()));
|
||||
}
|
||||
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
|
||||
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
|
||||
return new KafkaConsumer<>(props);
|
||||
|
|
@ -36,6 +39,9 @@ public class KafkaConsumerFactory {
|
|||
}
|
||||
props.put("enable.auto.commit", kafkaSettings.getEnableAutoCommit().toString());
|
||||
props.put("session.timeout.ms", kafkaSettings.getSessionTimeoutMs().toString());
|
||||
if (kafkaSettings.getMaxPollIntervalMs() != null) {
|
||||
props.put("max.poll.interval.ms", String.valueOf(kafkaSettings.getMaxPollIntervalMs()));
|
||||
}
|
||||
props.put("auto.offset.reset", kafkaSettings.getAutoOffsetReset());
|
||||
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
|
||||
// props.put("value.deserializer", JsonDeserializer.class.getName());
|
||||
|
|
|
|||
|
|
@ -6,6 +6,7 @@ public class KafkaConsumerSettings {
|
|||
private Boolean enableAutoCommit;
|
||||
private Integer sessionTimeoutMs;
|
||||
private String autoOffsetReset;
|
||||
private Integer maxPollIntervalMs = null;
|
||||
|
||||
|
||||
public String getBootstrapServers() {
|
||||
|
|
@ -47,4 +48,13 @@ public class KafkaConsumerSettings {
|
|||
public void setAutoOffsetReset(String autoOffsetReset) {
|
||||
this.autoOffsetReset = autoOffsetReset;
|
||||
}
|
||||
|
||||
|
||||
public Integer getMaxPollIntervalMs() {
|
||||
return maxPollIntervalMs;
|
||||
}
|
||||
|
||||
public void setManuallyMaxPollIntervalMs(Integer maxPollIntervalMs) {
|
||||
this.maxPollIntervalMs = maxPollIntervalMs;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue