clearing-service changeStatusExtract => gateway without timeout

optional max-poll-interval-ms setting
This commit is contained in:
ialbert 2024-01-17 13:22:13 +03:00
parent ab80d194a5
commit 0fd9208e41
6 changed files with 61 additions and 29 deletions

View file

@ -2,6 +2,8 @@ package ru.spcex.clearing.config;
import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.producer.Producer; 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.annotation.Autowired;
import org.springframework.beans.factory.config.ConfigurableBeanFactory; import org.springframework.beans.factory.config.ConfigurableBeanFactory;
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Bean;
@ -26,19 +28,31 @@ import java.util.function.Supplier;
@Configuration @Configuration
public class KafkaConfig { public class KafkaConfig {
@Autowired private final Logger log = LoggerFactory.getLogger(getClass());
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) private final KafkaConsumerSettings kafkaSettings;
@Bean("kafkaConsumer")
public Consumer<String, Object> createConsumer(ClearingServiceSettings settings) { public KafkaConfig(ClearingServiceSettings settings) {
return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); 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") @Bean("kafkaConsumerGateway")
public Supplier<Consumer<String, Object>> getwaySessionConsumer(ClearingServiceSettings settings) { public Supplier<Consumer<String, Object>> getwaySessionConsumer() {
AtomicInteger groupId = new AtomicInteger(1); AtomicInteger groupId = new AtomicInteger(1);
return () -> { return () -> {
KafkaConsumerSettings consumer = settings.getKafkaConsumer(); KafkaConsumerSettings consumer = kafkaSettings;
KafkaConsumerSettings gatewayConsumer = new KafkaConsumerSettings(); KafkaConsumerSettings gatewayConsumer = new KafkaConsumerSettings();
gatewayConsumer.setBootstrapServers(consumer.getBootstrapServers()); gatewayConsumer.setBootstrapServers(consumer.getBootstrapServers());
gatewayConsumer.setSessionTimeoutMs(consumer.getSessionTimeoutMs()); gatewayConsumer.setSessionTimeoutMs(consumer.getSessionTimeoutMs());

View file

@ -368,6 +368,7 @@ public class ValidationConfig {
new PresentById(IMDGDistributedNames.Map_Registry, ClearingError.RecordNotFound, true), new PresentById(IMDGDistributedNames.Map_Registry, ClearingError.RecordNotFound, true),
StatusExtractValidationRule.RegistryCodeCheck, StatusExtractValidationRule.RegistryCodeCheck,
StatusExtractValidationRule.ContractCheck, StatusExtractValidationRule.ContractCheck,
StatusExtractValidationRule.RequestStatusValid,
StatusExtractValidationRule.RegistryStatusValid StatusExtractValidationRule.RegistryStatusValid
); );
}; };

View file

@ -48,16 +48,11 @@ 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.Collection; import java.util.*;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.function.Function; import java.util.function.Function;
import java.util.function.Supplier;
import java.util.stream.Collectors; import java.util.stream.Collectors;
import java.util.stream.Stream; import java.util.stream.Stream;
import static ru.spcex.clearing.session.stage.impl.GatewayRequester.mapError;
import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD; import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD;
@Service @Service
@ -290,23 +285,11 @@ public class RegistryService {
return reqHelp.error(req.getId(), err.get()); return reqHelp.error(req.getId(), err.get());
} }
Registry rgs = validator.getStored(Stored.PresentById); Registry rgs = validator.getStored(Stored.PresentById);
Supplier<AssetOperationRequest> gtwBuilder = () -> GatewayRequestCreator.from(rgs, InOutDirection.in);
Optional<Boolean> gatewayOk;
if (trdTime.isTradingTime()) { 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 { } else {
log.debug("RegistryChangeStatusExtractRequest rgs.id={} not sending gateway request", rgs.getId()); 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()); rgs.setRegistryStatus(payload.getRegistryStatus());
log.debug("changing registry.id={} status to {}", rgs.getId(), 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()); 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) { public RequestInfoUpdate identificationFunds(BaseRequest<IdentificationFundsRequest> req) {
if (!rights.userHasRole(req.getUserId(), UserRole.Admin)) { if (!rights.userHasRole(req.getUserId(), UserRole.Admin)) {
return reqHelp.error(req.getId(), ClearingError.UserVerifyDenial); return reqHelp.error(req.getId(), ClearingError.UserVerifyDenial);

View file

@ -63,7 +63,7 @@ public enum StatusExtractValidationRule implements IValidationRule<ImdgValidatio
return empty(); return empty();
} }
}, },
RegistryStatusValid() { RequestStatusValid() {
@Override @Override
public Optional<EnumMessage> validate(ImdgValidationContext<RegistryChangeStatusExtractRequest> context) { public Optional<EnumMessage> validate(ImdgValidationContext<RegistryChangeStatusExtractRequest> context) {
RegistryChangeStatusExtractRequest validatedObject = context.getValidatedObject(); RegistryChangeStatusExtractRequest validatedObject = context.getValidatedObject();
@ -76,6 +76,17 @@ public enum StatusExtractValidationRule implements IValidationRule<ImdgValidatio
} }
return empty(); 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); private final static Logger log = LoggerFactory.getLogger(StatusExtractValidationRule.class);

View file

@ -22,6 +22,9 @@ public class KafkaConsumerFactory {
props.put("enable.auto.commit", kafkaSettings.getEnableAutoCommit().toString()); props.put("enable.auto.commit", kafkaSettings.getEnableAutoCommit().toString());
props.put("session.timeout.ms", kafkaSettings.getSessionTimeoutMs().toString()); props.put("session.timeout.ms", kafkaSettings.getSessionTimeoutMs().toString());
props.put("auto.offset.reset", kafkaSettings.getAutoOffsetReset()); 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("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
return new KafkaConsumer<>(props); return new KafkaConsumer<>(props);
@ -36,6 +39,9 @@ public class KafkaConsumerFactory {
} }
props.put("enable.auto.commit", kafkaSettings.getEnableAutoCommit().toString()); props.put("enable.auto.commit", kafkaSettings.getEnableAutoCommit().toString());
props.put("session.timeout.ms", kafkaSettings.getSessionTimeoutMs().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("auto.offset.reset", kafkaSettings.getAutoOffsetReset());
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
// props.put("value.deserializer", JsonDeserializer.class.getName()); // props.put("value.deserializer", JsonDeserializer.class.getName());

View file

@ -6,6 +6,7 @@ public class KafkaConsumerSettings {
private Boolean enableAutoCommit; private Boolean enableAutoCommit;
private Integer sessionTimeoutMs; private Integer sessionTimeoutMs;
private String autoOffsetReset; private String autoOffsetReset;
private Integer maxPollIntervalMs = null;
public String getBootstrapServers() { public String getBootstrapServers() {
@ -47,4 +48,13 @@ public class KafkaConsumerSettings {
public void setAutoOffsetReset(String autoOffsetReset) { public void setAutoOffsetReset(String autoOffsetReset) {
this.autoOffsetReset = autoOffsetReset; this.autoOffsetReset = autoOffsetReset;
} }
public Integer getMaxPollIntervalMs() {
return maxPollIntervalMs;
}
public void setManuallyMaxPollIntervalMs(Integer maxPollIntervalMs) {
this.maxPollIntervalMs = maxPollIntervalMs;
}
} }