xml-importer minor bug fix and refactoring

This commit is contained in:
Ivan Nikolaev-Axenov 2024-09-13 11:37:53 +03:00
parent b3420edb46
commit fc3fc9cef6
10 changed files with 31 additions and 40 deletions

View file

@ -32,8 +32,8 @@ public class ImporterImdgConfig {
}
@Autowired
@Bean
public HazelcastService imdgProvider(
@Bean(name = "hazelcastService")
public HazelcastService hazelcastService(
@Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer,
@Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter,
ImportXMLServiceSettings settings

View file

@ -3,6 +3,7 @@ package ru.spcex.clearing.xml.importer.config;
import java.util.function.Supplier;
import org.apache.kafka.clients.consumer.Consumer;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@ -15,26 +16,26 @@ import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.xml.importer.config.settings.ImportXMLServiceSettings;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
@Configuration
public class KafkaConfig {
@Bean
public ProducerFactory<String, Object> pf(ImportXMLServiceSettings settings) {
@Bean(name = "producerFactory")
public ProducerFactory<String, Object> producerFactory(ImportXMLServiceSettings settings) {
KafkaProducerSettings kafkaSettings = settings.getKafkaProducer();
return KafkaProducerFactory.producerFactory(kafkaSettings);
}
@Bean
public KafkaTemplate<String, Object> kafkaTemplate(ProducerFactory<String, Object> pf) {
@Bean(name = "kafkaTemplate")
public KafkaTemplate<String, Object> kafkaTemplate(@Qualifier("producerFactory") ProducerFactory<String, Object> pf) {
return new KafkaTemplate<>(pf);
}
@Autowired
@Bean
public Supplier<KafkaSender> kafkaSender(KafkaTemplate<String, Object> kafkaTemplate,
ImdgProvider imdgProvider) {
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
@Bean(name = "kafkaSender")
public Supplier<KafkaSender> kafkaSender(@Qualifier("kafkaTemplate") KafkaTemplate<String, Object> kafkaTemplate,
@Qualifier("hazelcastService") HazelcastService hazelcastService) {
ImdgId imdgIdGenerator = hazelcastService.getImdgIdGenerator();
return () -> KafkaSender
.setup()
.setKafkaTemplate(kafkaTemplate)
@ -45,7 +46,7 @@ public class KafkaConfig {
@Autowired
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
@Bean
@Bean(name = "createConsumer")
public Consumer<String, Object> createConsumer(ImportXMLServiceSettings settings) {
return KafkaConsumerFactory.consumer(settings.getKafkaConsumer());
}

View file

@ -99,7 +99,7 @@ public class SFTPConfig {
@ConditionalOnProperty(value = "import-xml-service.process-sdf-files", havingValue = "true")
@ServiceActivator(inputChannel = "listSftpChannelSdf")
public MessageHandler handlerListSdf(@Qualifier("sftpSessionFactorySdf") SessionFactory<ChannelSftp.LsEntry> sessionFactory,
ImportXMLServiceSettings settings) {
ImportXMLServiceSettings settings) {
SftpOutboundGateway sftpOutboundGateway = new SftpOutboundGateway(sessionFactory, MGET.getCommand(), null);
sftpOutboundGateway.setLocalDirectory(new File(settings.getStoreSdf().getSrcDir()));
sftpOutboundGateway.setAutoCreateLocalDirectory(true);

View file

@ -13,9 +13,9 @@ import javax.xml.stream.XMLInputFactory;
import javax.xml.stream.XMLOutputFactory;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.context.ApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
@ -29,22 +29,19 @@ import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.xml.importer.config.settings.ImportXMLServiceSettings;
import ru.spcex.clearing.xml.importer.logic.data.enums.ETable;
import ru.spcex.platform.classes.base.SpcexObjectBase;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.iml.hazelcast.adapter.ImdgHazelcast;
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
@Configuration
@EnableConfigurationProperties
public class XMLImporterConfig {
private final Logger log = LoggerFactory.getLogger(getClass());
private final ImportXMLServiceSettings settings;
private final ApplicationContext context;
private final ImdgProvider hazelcastService;
private final HazelcastService hazelcastService;
public XMLImporterConfig(ImportXMLServiceSettings settings,
ApplicationContext context,
ImdgProvider hazelcastService) {
@Qualifier("hazelcastService") HazelcastService hazelcastService) {
this.settings = settings;
this.context = context;
this.hazelcastService = hazelcastService;
}

View file

@ -2,6 +2,7 @@ package ru.spcex.clearing.xml.importer.logic.steps;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.stereotype.Component;
import ru.clearing.classes.statics.data.sdf.SDf57;
import ru.spcex.clearing.platform.messaging.domain.cud.account.BankAccountNewRequest;
@ -25,7 +26,7 @@ public class ImportToDB {
private final HazelcastService hazelcastService;
private final XmlImportKafkaMessenger kafkaMessenger;
public ImportToDB(HazelcastService hazelcastService,
public ImportToDB(@Qualifier("hazelcastService") HazelcastService hazelcastService,
XmlImportKafkaMessenger kafkaMessenger) {
this.hazelcastService = hazelcastService;
this.kafkaMessenger = kafkaMessenger;

View file

@ -2,13 +2,14 @@ package ru.spcex.clearing.xml.importer.logic.steps;
import static ru.spcex.clearing.platform.messaging.domain.Consts.PAIR_SDF;
import java.util.HashMap;
import java.util.EnumMap;
import java.util.Map;
import java.util.function.Consumer;
import java.util.function.Supplier;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.stereotype.Component;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.account.BankAccountNewRequest;
@ -18,7 +19,6 @@ import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationNew
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.PairSdfRequest;
import ru.spcex.clearing.platform.messaging.serialization.LogFormatter;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.xml.importer.config.settings.ImportXMLServiceSettings;
import ru.spcex.clearing.xml.importer.logic.data.ResultContainer;
import ru.spcex.clearing.xml.importer.logic.data.enums.ETable;
import ru.spcex.clearing.xml.importer.logic.data.enums.StageResult;
@ -30,15 +30,10 @@ import ru.spcex.platform.enumeration.SdfTable;
public class XmlImportKafkaMessenger implements InitializingBean {
final Logger log = LoggerFactory.getLogger(getClass());
private final Supplier<KafkaSender> kafka;
private final Map<ETable, Consumer<Long>> messengers;
private final Map<ETable, Consumer<Long>> messengers = new EnumMap<>(ETable.class);
private final ImportXMLServiceSettings settings;
public XmlImportKafkaMessenger(Supplier<KafkaSender> kafka,
ImportXMLServiceSettings settings) {
public XmlImportKafkaMessenger(@Qualifier("kafkaSender") Supplier<KafkaSender> kafka) {
this.kafka = kafka;
this.settings = settings;
this.messengers = new HashMap<>();
}
@Override

View file

@ -57,6 +57,7 @@ public class FileChecker {
List<File> xmlFiles = srcDir.stream()
.map(this::lsXML)
.flatMap(List::stream)
.distinct()
.toList();
if (xmlFiles.isEmpty()) return newFiles;

View file

@ -64,8 +64,8 @@ public class ImporterImdgTestConfig {
}
@Autowired
@Bean(name = "imdgProvider")
public HazelcastService imdgTestProvider(
@Bean(name = "hazelcastService")
public HazelcastService hazelcastServiceTest(
@Qualifier("taskExecutorHazelcastTestClientInitializerXmlImporter") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer,
@Qualifier("taskExecutorTestIdGeneratorAwaiterXmlImporter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter,
@Qualifier("hazelcastClientParamsXmlImporter") HazelcastClientParams params) {

View file

@ -17,7 +17,6 @@ import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.mockito.ArgumentCaptor;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
import org.springframework.boot.test.mock.mockito.MockReset;
@ -30,13 +29,10 @@ import org.springframework.kafka.requestreply.ReplyingKafkaTemplate;
import org.springframework.kafka.requestreply.RequestReplyFuture;
import org.springframework.kafka.support.SendResult;
import org.springframework.util.concurrent.ListenableFuture;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.test.TestUtils;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
@Configuration
@Import(ImporterImdgTestConfig.class)
@ -132,9 +128,9 @@ public class KafkaTestConfig {
@Bean("kafkaSender")
public Supplier<KafkaSender> kafkaSender(@Qualifier("kafkaTestTemplate") KafkaTemplate<String, Object> kafkaTemplate,
@Qualifier("imdgProvider") ImdgProvider imdgProvider,
@Qualifier("hazelcastService") HazelcastService hazelcastService,
Producer<String, Object> mockProducer) {
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
ImdgId imdgIdGenerator = hazelcastService.getImdgIdGenerator();
return () -> KafkaSender
.setup()
.setKafkaTemplate(kafkaTemplate)

View file

@ -59,7 +59,7 @@ class ImportToDBTest {
private final DateTimeFormatter formatter = DateTimeFormatter.ofPattern("dd.MM.yy");
@Autowired
@Qualifier("imdgProvider")
@Qualifier("hazelcastService")
private HazelcastService hazelcastService;
@Autowired